From cb115bd6d5bf70c146f008a30fe3c2a611604e84 Mon Sep 17 00:00:00 2001 From: Dylan Myers Date: Thu, 6 Aug 2026 09:15:29 -0400 Subject: [PATCH] feat(o11y): self-tracing via runtime spans, adapter emit spans, and OTLP trace export (PIPE-1066) Assisted-by: Claude Opus 4.8 --- cmd/blitz/main.go | 34 ++++++++--- embed/new.go | 2 +- embed/telemetry.go | 12 ++++ embed/telemetry_test.go | 25 ++++++++ go.mod | 3 + go.sum | 6 ++ internal/config/config.go | 8 ++- internal/config/override.go | 3 + internal/config/override_test.go | 20 ++++++ internal/config/telemetry.go | 34 +++++++++++ internal/runtime/runtime.go | 47 ++++++++++++++- internal/runtime/runtime_span_test.go | 48 +++++++++++++++ internal/runtime/runtime_test.go | 10 +-- internal/service/service.go | 4 +- internal/service/service_test.go | 14 ++--- internal/telemetry/traces/traces.go | 77 ++++++++++++++++++++++++ internal/telemetry/traces/traces_test.go | 58 ++++++++++++++++++ output/adapter.go | 48 +++++++++++---- output/adapter_test.go | 64 +++++++++++++++++--- package/completions/blitz.bash | 12 ++++ 20 files changed, 481 insertions(+), 48 deletions(-) create mode 100644 internal/config/telemetry.go create mode 100644 internal/runtime/runtime_span_test.go create mode 100644 internal/telemetry/traces/traces.go create mode 100644 internal/telemetry/traces/traces_test.go diff --git a/cmd/blitz/main.go b/cmd/blitz/main.go index 1f45c61..1c97659 100644 --- a/cmd/blitz/main.go +++ b/cmd/blitz/main.go @@ -25,6 +25,7 @@ import ( "github.com/observiq/blitz/internal/logging" "github.com/observiq/blitz/internal/service" "github.com/observiq/blitz/internal/telemetry/metrics" + "github.com/observiq/blitz/internal/telemetry/traces" "github.com/observiq/blitz/output" fileout "github.com/observiq/blitz/output/file" hecout "github.com/observiq/blitz/output/hec" @@ -143,10 +144,27 @@ func run(cmd *cobra.Command, args []string) error { cancel() }() - // Blitz routes its own self-telemetry through this bundle. Standalone - // leaves the providers nil so they fall back to the process-global - // provider configured by setupMetrics (Prometheus). - tel := embed.TelemetrySettings{Logger: logger} + // Blitz routes its own self-telemetry through this bundle. Metrics leave + // the provider nil so they fall back to the process-global provider + // configured by setupMetrics (Prometheus). Trace export is opt-in via the + // telemetry.traces config: when an OTLP endpoint is set, spans export + // there; otherwise the nil TracerProvider means spans are created but + // dropped by the global no-op provider. + tel := embed.TelemetrySettings{ + Logger: logger, + PerBatchSpans: cfg.Telemetry.Traces.PerBatchSpans, + } + if cfg.Telemetry.Traces.OTLPEndpoint != "" { + otlpTraces, terr := traces.NewOTLP(ctx, cfg.Telemetry.Traces.OTLPEndpoint, cfg.Telemetry.Traces.Insecure) + if terr != nil { + logger.Error("Failed to enable self-telemetry trace export", zap.Error(terr)) + return terr + } + defer func() { _ = otlpTraces.Shutdown(context.Background()) }() + tel.TracerProvider = otlpTraces.Provider() + logger.Info("self-telemetry trace export enabled", + zap.String("endpoint", cfg.Telemetry.Traces.OTLPEndpoint)) + } // Configure output first var outputInstance output.Output @@ -355,7 +373,7 @@ func run(cmd *cobra.Command, args []string) error { // Set up SIGUSR1 restart signal handler setupRestartSignal(ctx, logger, tracker) - svc, err := service.New(logger, generators, outputInstance) + svc, err := service.New(logger, generators, outputInstance, tel) if err != nil { logger.Error("Failed to create service", zap.Error(err)) return err @@ -427,13 +445,13 @@ func createGenerator(logger *zap.Logger, genCfg config.Generator, out output.Out // when an output doesn't support a signal the configured generator // needs. consumers := dispatch.EmbedConsumers{ - LogConsumer: output.WriterAsLogConsumer(out), + LogConsumer: output.WriterAsLogConsumer(out, tel), } if mw, ok := out.(output.MetricWriter); ok { - consumers.MetricConsumer = output.WriterAsMetricConsumer(mw) + consumers.MetricConsumer = output.WriterAsMetricConsumer(mw, tel) } if tw, ok := out.(output.TraceWriter); ok { - consumers.TraceConsumer = output.WriterAsTraceConsumer(tw) + consumers.TraceConsumer = output.WriterAsTraceConsumer(tw, tel) } mod, err := dispatch.ForEmbed(logger, genCfg, consumers, nil, tel) if err != nil { diff --git a/embed/new.go b/embed/new.go index db76cc2..802046b 100644 --- a/embed/new.go +++ b/embed/new.go @@ -54,7 +54,7 @@ func (r *runner) Start(ctx context.Context, host Host) error { for i, m := range r.cfg.Modules { rtModules[i] = m } - rt := runtime.New(logger, rtModules) + rt := runtime.New(logger, rtModules, host.TracerProvider) if err := rt.Start(ctx); err != nil { return err } diff --git a/embed/telemetry.go b/embed/telemetry.go index 90889aa..29b887a 100644 --- a/embed/telemetry.go +++ b/embed/telemetry.go @@ -1,6 +1,7 @@ package embed import ( + "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/metric" metricnoop "go.opentelemetry.io/otel/metric/noop" "go.opentelemetry.io/otel/trace" @@ -35,6 +36,17 @@ type TelemetrySettings struct { PerBatchSpans bool } +// Tracer returns a tracer for the given instrumentation scope from the bundle's +// TracerProvider. A nil TracerProvider falls back to the process global, so the +// result is always safe to use. +func (t TelemetrySettings) Tracer(scope string) trace.Tracer { + tp := t.TracerProvider + if tp == nil { + tp = otel.GetTracerProvider() + } + return tp.Tracer(scope) +} + // NopTelemetry returns a TelemetrySettings wired to no-op providers and a nop // logger. Use it where a caller has no telemetry to route, most commonly in // tests and in construction paths that record nothing. It is distinct from a diff --git a/embed/telemetry_test.go b/embed/telemetry_test.go index fcbdfe3..562a6f5 100644 --- a/embed/telemetry_test.go +++ b/embed/telemetry_test.go @@ -1,9 +1,12 @@ package embed import ( + "context" "testing" "github.com/stretchr/testify/require" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" ) func TestNopTelemetry(t *testing.T) { @@ -28,3 +31,25 @@ func TestTelemetrySettings_zeroValueFieldsAreNil(t *testing.T) { require.Nil(t, tel.MeterProvider) require.Nil(t, tel.TracerProvider) } + +func TestTelemetrySettings_Tracer_usesProvidedProvider(t *testing.T) { + exporter := tracetest.NewInMemoryExporter() + tp := sdktrace.NewTracerProvider(sdktrace.WithSyncer(exporter)) + tel := TelemetrySettings{TracerProvider: tp} + + _, span := tel.Tracer("test").Start(context.Background(), "op") + span.End() + + spans := exporter.GetSpans() + require.Len(t, spans, 1) + require.Equal(t, "op", spans[0].Name) +} + +func TestTelemetrySettings_Tracer_nilFallsBackToGlobal(t *testing.T) { + var tel TelemetrySettings + + require.NotPanics(t, func() { + _, span := tel.Tracer("test").Start(context.Background(), "op") + span.End() + }) +} diff --git a/go.mod b/go.mod index 8d7e540..8137e23 100644 --- a/go.mod +++ b/go.mod @@ -20,6 +20,8 @@ require ( github.com/spf13/viper v1.21.0 github.com/stretchr/testify v1.11.1 go.opentelemetry.io/otel v1.44.0 + go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.44.0 + go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.44.0 go.opentelemetry.io/otel/exporters/prometheus v0.66.0 go.opentelemetry.io/otel/metric v1.44.0 go.opentelemetry.io/otel/sdk v1.44.0 @@ -42,6 +44,7 @@ require ( github.com/anthropics/anthropic-sdk-go v1.19.0 // indirect github.com/beorn7/perks v1.0.1 // indirect github.com/ccojocar/zxcvbn-go v1.0.4 // indirect + github.com/cenkalti/backoff/v5 v5.0.3 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect github.com/cpuguy83/go-md2man/v2 v2.0.7 // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect diff --git a/go.sum b/go.sum index 41c3701..4e4a2c8 100644 --- a/go.sum +++ b/go.sum @@ -18,6 +18,8 @@ github.com/ccojocar/zxcvbn-go v1.0.4 h1:FWnCIRMXPj43ukfX000kvBZvV6raSxakYr1nzyNr github.com/ccojocar/zxcvbn-go v1.0.4/go.mod h1:3GxGX+rHmueTUMvm5ium7irpyjmm7ikxYFOSJB21Das= github.com/cenkalti/backoff/v4 v4.3.0 h1:MyRJ/UdXutAwSAT+s3wNd7MfTIcy71VQueUuFK343L8= github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyYozVcomhLiZE= +github.com/cenkalti/backoff/v5 v5.0.3 h1:ZN+IMa753KfX5hd8vVaMixjnqRZ3y8CuJKRKj1xcsSM= +github.com/cenkalti/backoff/v5 v5.0.3/go.mod h1:rkhZdG3JZukswDf7f0cwqPNk4K0sa+F97BxZthm/crw= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/cpuguy83/go-md2man/v2 v2.0.6/go.mod h1:oOW0eioCTA6cOiMLiUPZOpcVxMig6NIQQ7OS05n1F4g= @@ -159,6 +161,10 @@ go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0 h1:F7Jx+6h go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.61.0/go.mod h1:UHB22Z8QsdRDrnAtX4PntOl36ajSxcdUMt1sF7Y6E7Q= go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU= go.opentelemetry.io/otel v1.44.0/go.mod h1:BMgjTHL9WPRlRjL2oZCBTL4whCGtXch2H4BhOPIAyYc= +go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.44.0 h1:4YsVu3B8+3qtWYYrsUYgn0OG78pN0rnNPRGX4SbokQI= +go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.44.0/go.mod h1:+wnlSn0mD1ADVMe3v9Z/WIaiz6q6gL2J/ejaAmdmv80= +go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.44.0 h1:qazEJlUOQzhCpzQpFETGby7EdqjI1wsd0W+6Gg1SCTU= +go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc v1.44.0/go.mod h1:fOD2Yefuxixkx3ahVNf0O/PERb6r4OlbxfATVnYvzCo= go.opentelemetry.io/otel/exporters/prometheus v0.66.0 h1:vkrK8PAznv2NKt2r+kdu252ccGzkEqLc2aSXbQIALYQ= go.opentelemetry.io/otel/exporters/prometheus v0.66.0/go.mod h1:V/UB6D3vMF/UBOL5igAsAYnk1nG/bzYYTzvsB16cy7o= go.opentelemetry.io/otel/metric v1.44.0 h1:1w0gILTcHdr3YI+ixLyjemwrVnsMURbTZFrSYCdDdmc= diff --git a/internal/config/config.go b/internal/config/config.go index b5ad4d3..c9e3188 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -17,8 +17,11 @@ type Config struct { Generators []Generator `yaml:"generators,omitempty" mapstructure:"generators,omitempty"` // Output configuration Output Output `yaml:"output,omitempty" mapstructure:"output,omitempty"` - // Metrics configuration + // Metrics configuration (Prometheus scrape endpoint for self-metrics) Metrics Metrics `yaml:"metrics,omitempty" mapstructure:"metrics,omitempty"` + // Telemetry configures export of blitz's own self-telemetry (self-traces, + // and later self-logs) via OTLP. + Telemetry Telemetry `yaml:"telemetry,omitempty" mapstructure:"telemetry,omitempty"` // OnFinish controls behavior when finite generation completes. // One of: "exit" (default), "idle" OnFinish string `yaml:"onFinish,omitempty" mapstructure:"onFinish,omitempty"` @@ -38,6 +41,9 @@ func (c *Config) Validate() error { if err := c.Metrics.Validate(); err != nil { return err } + if err := c.Telemetry.Validate(); err != nil { + return err + } if c.OnFinish != "" && c.OnFinish != "exit" && c.OnFinish != "idle" { return fmt.Errorf("onFinish must be one of: exit, idle, got %q", c.OnFinish) } diff --git a/internal/config/override.go b/internal/config/override.go index 5c57781..1153206 100644 --- a/internal/config/override.go +++ b/internal/config/override.go @@ -338,6 +338,9 @@ func DefaultOverrides() []*Override { NewOverride("output.hec.source", "HEC event source metadata", DefaultHECSource), NewOverride("output.hec.sourceType", "HEC event sourcetype metadata", DefaultHECSourceType), NewOverride("output.hec.index", "HEC target index (empty = token default)", ""), + NewOverride("telemetry.traces.otlpEndpoint", "OTLP gRPC endpoint (host:port) for exporting blitz's own spans (empty = disabled)", ""), + NewOverride("telemetry.traces.insecure", "send blitz's own spans over plaintext gRPC (no TLS)", false), + NewOverride("telemetry.traces.perBatchSpans", "enable higher-volume per-emit-cycle spans (off by default)", false), } overrides = append(overrides, tcpTLSOverrides()...) diff --git a/internal/config/override_test.go b/internal/config/override_test.go index afde3b0..3907a16 100644 --- a/internal/config/override_test.go +++ b/internal/config/override_test.go @@ -140,6 +140,9 @@ func getTestOverrideFlagsArgs() []string { "--output-hec-tls-min-version", "1.3", "--output-stdout-flushinterval", "50ms", "--metrics-port", "8080", + "--telemetry-traces-otlpendpoint", "traces.example:4317", + "--telemetry-traces-insecure", "true", + "--telemetry-traces-perbatchspans", "true", } } @@ -272,6 +275,9 @@ func getTestOverrideEnvs() map[string]string { "BLITZ_OUTPUT_HEC_TLS_MIN_VERSION": "1.2", "BLITZ_OUTPUT_STDOUT_FLUSHINTERVAL": "75ms", "BLITZ_METRICS_PORT": "9100", + "BLITZ_TELEMETRY_TRACES_OTLPENDPOINT": "traces.env.example:4317", + "BLITZ_TELEMETRY_TRACES_INSECURE": "true", + "BLITZ_TELEMETRY_TRACES_PERBATCHSPANS": "true", } } @@ -677,6 +683,13 @@ func TestOverrideFlags(t *testing.T) { Metrics: Metrics{ Port: 8080, }, + Telemetry: Telemetry{ + Traces: TracesTelemetry{ + OTLPEndpoint: "traces.example:4317", + Insecure: true, + PerBatchSpans: true, + }, + }, } require.Equal(t, expectedCfg, cfg) } @@ -886,6 +899,13 @@ func TestOverrideEnvs(t *testing.T) { Metrics: Metrics{ Port: 9100, }, + Telemetry: Telemetry{ + Traces: TracesTelemetry{ + OTLPEndpoint: "traces.env.example:4317", + Insecure: true, + PerBatchSpans: true, + }, + }, } require.Equal(t, expectedCfg, cfg) } diff --git a/internal/config/telemetry.go b/internal/config/telemetry.go new file mode 100644 index 0000000..b1796ff --- /dev/null +++ b/internal/config/telemetry.go @@ -0,0 +1,34 @@ +package config + +// Telemetry configures export of blitz's OWN self-telemetry (its internal +// logs, metrics, and traces). This is distinct from the data blitz generates. +// The existing `metrics` block still governs the Prometheus scrape endpoint +// for self-metrics; this block adds OTLP export for self-traces (and, in a +// later phase, self-logs). +type Telemetry struct { + // Traces configures OTLP export of blitz's internal spans. + Traces TracesTelemetry `yaml:"traces,omitempty" mapstructure:"traces,omitempty"` +} + +// TracesTelemetry configures OTLP gRPC export of blitz's internal spans. +type TracesTelemetry struct { + // OTLPEndpoint is the OTLP gRPC endpoint (host:port). Empty disables trace + // export: spans are still created but routed to a no-op provider. + OTLPEndpoint string `yaml:"otlpEndpoint,omitempty" mapstructure:"otlpEndpoint,omitempty"` + + // Insecure sends spans over plaintext gRPC (no TLS). Defaults to false. + Insecure bool `yaml:"insecure,omitempty" mapstructure:"insecure,omitempty"` + + // PerBatchSpans enables the higher-volume per-emit-cycle spans. Off by + // default; the coarse session and generator-lifecycle spans do not depend + // on it. + PerBatchSpans bool `yaml:"perBatchSpans,omitempty" mapstructure:"perBatchSpans,omitempty"` +} + +// Validate validates the telemetry configuration. Export is off by default and +// all fields are optional, so there is nothing to reject today; the method +// exists to match the config-block Validate convention and to host future +// checks (e.g. endpoint format). +func (t Telemetry) Validate() error { + return nil +} diff --git a/internal/runtime/runtime.go b/internal/runtime/runtime.go index aedd778..23f901e 100644 --- a/internal/runtime/runtime.go +++ b/internal/runtime/runtime.go @@ -4,9 +4,16 @@ import ( "context" "fmt" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/trace" "go.uber.org/zap" ) +// tracerScope is the instrumentation scope for the runtime's self-telemetry +// spans. +const tracerScope = "github.com/observiq/blitz/internal/runtime" + // Module is the narrow lifecycle interface Runtime operates on. The // embed.ProducerModule type can be wrapped to satisfy this contract — // Runtime stays decoupled from the embed package to avoid an import @@ -31,19 +38,38 @@ type Module interface { // host-level concerns. CLI adds signal handling, YAML loading, and // output wiring; embed adds host-supplied consumers and resource // attributes. +// +// Runtime emits blitz's session-level self-telemetry: a root "blitz.session" +// span covering Start to Stop, with a child "blitz.generator.run" span per +// module (bounded by that module's lifetime). These spans are decoupled from +// the embed package; the caller passes a raw trace.TracerProvider. type Runtime struct { logger *zap.Logger modules []Module + tracer trace.Tracer + + // sessionSpan and moduleSpans hold the open self-telemetry spans between + // Start and Stop. Start and Stop are called once each and never + // concurrently, so no synchronization is needed. moduleSpans is + // index-aligned with modules. + sessionSpan trace.Span + moduleSpans []trace.Span } -// New returns a Runtime configured with the given logger and modules. -func New(logger *zap.Logger, modules []Module) *Runtime { +// New returns a Runtime configured with the given logger, modules, and tracer +// provider. A nil tracerProvider falls back to the process global, so span +// emission is always safe. +func New(logger *zap.Logger, modules []Module, tracerProvider trace.TracerProvider) *Runtime { if logger == nil { logger = zap.NewNop() } + if tracerProvider == nil { + tracerProvider = otel.GetTracerProvider() + } return &Runtime{ logger: logger, modules: modules, + tracer: tracerProvider.Tracer(tracerScope), } } @@ -51,9 +77,15 @@ func New(logger *zap.Logger, modules []Module) *Runtime { // an error, Start stops the modules already started (in reverse order) // and returns the failure. func (r *Runtime) Start(ctx context.Context) error { + ctx, r.sessionSpan = r.tracer.Start(ctx, "blitz.session") + r.moduleSpans = make([]trace.Span, 0, len(r.modules)) + started := make([]Module, 0, len(r.modules)) for _, m := range r.modules { - if err := m.Start(ctx); err != nil { + mctx, mspan := r.tracer.Start(ctx, "blitz.generator.run", + trace.WithAttributes(attribute.String("blitz.generator.name", m.Name()))) + if err := m.Start(mctx); err != nil { + mspan.End() // Roll back: stop modules already started, in reverse order. for i := len(started) - 1; i >= 0; i-- { if stopErr := started[i].Stop(ctx); stopErr != nil { @@ -61,10 +93,13 @@ func (r *Runtime) Start(ctx context.Context) error { zap.String("module", started[i].Name()), zap.Error(stopErr)) } + r.moduleSpans[i].End() } + r.sessionSpan.End() return fmt.Errorf("start module %s: %w", m.Name(), err) } started = append(started, m) + r.moduleSpans = append(r.moduleSpans, mspan) } return nil } @@ -83,6 +118,12 @@ func (r *Runtime) Stop(ctx context.Context) error { firstErr = fmt.Errorf("stop module %s: %w", m.Name(), err) } } + if i < len(r.moduleSpans) && r.moduleSpans[i] != nil { + r.moduleSpans[i].End() + } + } + if r.sessionSpan != nil { + r.sessionSpan.End() } return firstErr } diff --git a/internal/runtime/runtime_span_test.go b/internal/runtime/runtime_span_test.go new file mode 100644 index 0000000..fc4ea53 --- /dev/null +++ b/internal/runtime/runtime_span_test.go @@ -0,0 +1,48 @@ +package runtime_test + +import ( + "context" + "testing" + + "github.com/observiq/blitz/internal/runtime" + "github.com/stretchr/testify/require" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" +) + +type spanTestModule struct{ n string } + +func (m spanTestModule) Name() string { return m.n } +func (m spanTestModule) Start(context.Context) error { return nil } +func (m spanTestModule) Stop(context.Context) error { return nil } + +// TestRuntime_emitsSessionAndGeneratorSpans confirms the runtime emits a single +// root session span plus one lifecycle span per module, tagged with the +// generator name, through the provided TracerProvider. +func TestRuntime_emitsSessionAndGeneratorSpans(t *testing.T) { + exp := tracetest.NewInMemoryExporter() + tp := sdktrace.NewTracerProvider(sdktrace.WithSyncer(exp)) + + rt := runtime.New(nil, []runtime.Module{ + spanTestModule{n: "json"}, + spanTestModule{n: "apache"}, + }, tp) + require.NoError(t, rt.Start(context.Background())) + require.NoError(t, rt.Stop(context.Background())) + + byName := map[string]int{} + genNames := map[string]int{} + for _, s := range exp.GetSpans() { + byName[s.Name]++ + for _, a := range s.Attributes { + if string(a.Key) == "blitz.generator.name" { + genNames[a.Value.AsString()]++ + } + } + } + + require.Equal(t, 1, byName["blitz.session"], "exactly one session span") + require.Equal(t, 2, byName["blitz.generator.run"], "one lifecycle span per module") + require.Equal(t, 1, genNames["json"]) + require.Equal(t, 1, genNames["apache"]) +} diff --git a/internal/runtime/runtime_test.go b/internal/runtime/runtime_test.go index a2d72ce..59bb803 100644 --- a/internal/runtime/runtime_test.go +++ b/internal/runtime/runtime_test.go @@ -46,7 +46,7 @@ func TestRuntime_StartCallsEveryModuleInOrder(t *testing.T) { b := &recordingModule{name: "b", startCall: startOrder} c := &recordingModule{name: "c", startCall: startOrder} - rt := runtime.New(zaptest.NewLogger(t), []runtime.Module{a, b, c}) + rt := runtime.New(zaptest.NewLogger(t), []runtime.Module{a, b, c}, nil) if err := rt.Start(context.Background()); err != nil { t.Fatalf("Start: %v", err) } @@ -67,7 +67,7 @@ func TestRuntime_StartRollsBackOnFailure(t *testing.T) { b := &recordingModule{name: "b", stopCall: stopOrder} failing := &recordingModule{name: "failing", startErr: errors.New("boom")} - rt := runtime.New(zaptest.NewLogger(t), []runtime.Module{a, b, failing}) + rt := runtime.New(zaptest.NewLogger(t), []runtime.Module{a, b, failing}, nil) err := rt.Start(context.Background()) if err == nil { t.Fatal("expected error from Start") @@ -95,7 +95,7 @@ func TestRuntime_StopCallsEveryModuleInReverseOrder(t *testing.T) { b := &recordingModule{name: "b", stopCall: stopOrder} c := &recordingModule{name: "c", stopCall: stopOrder} - rt := runtime.New(zaptest.NewLogger(t), []runtime.Module{a, b, c}) + rt := runtime.New(zaptest.NewLogger(t), []runtime.Module{a, b, c}, nil) if err := rt.Start(context.Background()); err != nil { t.Fatalf("Start: %v", err) } @@ -114,7 +114,7 @@ func TestRuntime_StopContinuesOnError(t *testing.T) { b := &recordingModule{name: "b", stopErr: errors.New("b-stop-fail")} c := &recordingModule{name: "c"} - rt := runtime.New(zaptest.NewLogger(t), []runtime.Module{a, b, c}) + rt := runtime.New(zaptest.NewLogger(t), []runtime.Module{a, b, c}, nil) if err := rt.Start(context.Background()); err != nil { t.Fatalf("Start: %v", err) } @@ -130,7 +130,7 @@ func TestRuntime_StopContinuesOnError(t *testing.T) { } func TestRuntime_NewWithNilLoggerUsesNop(t *testing.T) { - rt := runtime.New(nil, nil) + rt := runtime.New(nil, nil, nil) // Should not panic with empty modules. if err := rt.Start(context.Background()); err != nil { t.Errorf("Start with empty modules and nil logger: %v", err) diff --git a/internal/service/service.go b/internal/service/service.go index 40969ff..b5215e9 100644 --- a/internal/service/service.go +++ b/internal/service/service.go @@ -25,7 +25,7 @@ type Service struct { } // New creates a new service with multiple generators and a single output. -func New(logger *zap.Logger, generators []any, output output.Output) (*Service, error) { +func New(logger *zap.Logger, generators []any, output output.Output, tel embed.TelemetrySettings) (*Service, error) { if logger == nil { return nil, fmt.Errorf("logger cannot be nil") } @@ -54,7 +54,7 @@ func New(logger *zap.Logger, generators []any, output output.Output) (*Service, Logger: logger, Generators: generators, Output: output, - runtime: runtime.New(logger, modules), + runtime: runtime.New(logger, modules, tel.TracerProvider), legacy: legacy, }, nil } diff --git a/internal/service/service_test.go b/internal/service/service_test.go index be9fc8b..9003a68 100644 --- a/internal/service/service_test.go +++ b/internal/service/service_test.go @@ -48,34 +48,34 @@ func TestNew(t *testing.T) { logger := zaptest.NewLogger(t) t.Run("nil logger", func(t *testing.T) { - _, err := New(nil, []any{&stubLogGen{}}, &stubOutput{}) + _, err := New(nil, []any{&stubLogGen{}}, &stubOutput{}, embed.NopTelemetry()) require.Error(t, err) }) t.Run("nil generators", func(t *testing.T) { - _, err := New(logger, nil, &stubOutput{}) + _, err := New(logger, nil, &stubOutput{}, embed.NopTelemetry()) require.Error(t, err) }) t.Run("empty generators", func(t *testing.T) { - _, err := New(logger, []any{}, &stubOutput{}) + _, err := New(logger, []any{}, &stubOutput{}, embed.NopTelemetry()) require.Error(t, err) }) t.Run("nil output", func(t *testing.T) { - _, err := New(logger, []any{&stubLogGen{}}, nil) + _, err := New(logger, []any{&stubLogGen{}}, nil, embed.NopTelemetry()) require.Error(t, err) }) t.Run("valid single log generator", func(t *testing.T) { - svc, err := New(logger, []any{&stubLogGen{}}, &stubOutput{}) + svc, err := New(logger, []any{&stubLogGen{}}, &stubOutput{}, embed.NopTelemetry()) require.NoError(t, err) assert.NotNil(t, svc) }) t.Run("valid multi generator", func(t *testing.T) { gens := []any{&stubLogGen{}, &stubLogGen{}} - svc, err := New(logger, gens, &stubOutput{}) + svc, err := New(logger, gens, &stubOutput{}, embed.NopTelemetry()) require.NoError(t, err) assert.NotNil(t, svc) }) @@ -87,7 +87,7 @@ func TestStartStop(t *testing.T) { logGen := &stubLogGen{} logGen2 := &stubLogGen{} - svc, err := New(logger, []any{logGen, logGen2}, &stubOutput{}) + svc, err := New(logger, []any{logGen, logGen2}, &stubOutput{}, embed.NopTelemetry()) require.NoError(t, err) require.NoError(t, svc.Start()) diff --git a/internal/telemetry/traces/traces.go b/internal/telemetry/traces/traces.go new file mode 100644 index 0000000..4a7504a --- /dev/null +++ b/internal/telemetry/traces/traces.go @@ -0,0 +1,77 @@ +// Package traces provides an OTLP gRPC exporter and TracerProvider for +// blitz's own self-telemetry spans. It mirrors the Prometheus metrics +// setup in internal/telemetry/metrics: a thin wrapper around the OTel SDK +// that the standalone CLI wires up when trace export is configured. +package traces + +import ( + "context" + "fmt" + "os" + + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc" + "go.opentelemetry.io/otel/sdk/resource" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + semconv "go.opentelemetry.io/otel/semconv/v1.4.0" + "go.opentelemetry.io/otel/trace" +) + +const serviceName = "blitz" + +// osHostname and newTraceExporter are indirections over os.Hostname and the +// OTLP trace exporter constructor, so tests can exercise the hostname-fallback +// and exporter-error paths deterministically. +var ( + osHostname = os.Hostname + newTraceExporter = otlptracegrpc.New +) + +// OTLP owns an OTLP gRPC trace exporter and the TracerProvider built on it. +type OTLP struct { + provider *sdktrace.TracerProvider +} + +// NewOTLP builds a batching TracerProvider that exports blitz's self-telemetry +// spans over OTLP gRPC to endpoint (host:port). insecure sends over plaintext. +// It also installs the provider as the process global so any otel.Tracer user +// picks it up. The exporter connects lazily, so a nil error does not imply the +// collector is reachable. +func NewOTLP(ctx context.Context, endpoint string, insecure bool) (*OTLP, error) { + hostname, err := osHostname() + if err != nil { + hostname = "unknown" + } + res := resource.NewWithAttributes(semconv.SchemaURL, + semconv.ServiceNameKey.String(serviceName), + semconv.HostNameKey.String(hostname), + ) + + opts := []otlptracegrpc.Option{otlptracegrpc.WithEndpoint(endpoint)} + if insecure { + opts = append(opts, otlptracegrpc.WithInsecure()) + } + exporter, err := newTraceExporter(ctx, opts...) + if err != nil { + return nil, fmt.Errorf("create otlp trace exporter: %w", err) + } + + provider := sdktrace.NewTracerProvider( + sdktrace.WithBatcher(exporter), + sdktrace.WithResource(res), + ) + otel.SetTracerProvider(provider) + + return &OTLP{provider: provider}, nil +} + +// Provider returns the TracerProvider for wiring into a TelemetrySettings. +func (o *OTLP) Provider() trace.TracerProvider { return o.provider } + +// Shutdown flushes buffered spans and stops the exporter. +func (o *OTLP) Shutdown(ctx context.Context) error { + if o.provider != nil { + return o.provider.Shutdown(ctx) + } + return nil +} diff --git a/internal/telemetry/traces/traces_test.go b/internal/telemetry/traces/traces_test.go new file mode 100644 index 0000000..806d218 --- /dev/null +++ b/internal/telemetry/traces/traces_test.go @@ -0,0 +1,58 @@ +package traces + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/exporters/otlp/otlptrace" + "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc" +) + +// TestNewOTLP_buildsProviderAndShutsDown confirms the OTLP tracer provider +// constructs without a live collector (the exporter connects lazily) and shuts +// down cleanly. It does not assert export, which would require a collector. +func TestNewOTLP_buildsProviderAndShutsDown(t *testing.T) { + o, err := NewOTLP(context.Background(), "localhost:4317", true) + require.NoError(t, err) + require.NotNil(t, o) + require.NotNil(t, o.Provider()) + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + require.NoError(t, o.Shutdown(ctx)) +} + +// TestNewOTLP_hostnameFallback covers the os.Hostname error path: the provider +// still builds, falling back to an "unknown" hostname resource attribute. +func TestNewOTLP_hostnameFallback(t *testing.T) { + orig := osHostname + osHostname = func() (string, error) { return "", errors.New("no hostname") } + defer func() { osHostname = orig }() + + o, err := NewOTLP(context.Background(), "localhost:4317", true) + require.NoError(t, err) + require.NotNil(t, o) + require.NoError(t, o.Shutdown(context.Background())) +} + +// TestNewOTLP_exporterError covers the exporter-construction error path. +func TestNewOTLP_exporterError(t *testing.T) { + orig := newTraceExporter + newTraceExporter = func(context.Context, ...otlptracegrpc.Option) (*otlptrace.Exporter, error) { + return nil, errors.New("boom") + } + defer func() { newTraceExporter = orig }() + + _, err := NewOTLP(context.Background(), "localhost:4317", true) + require.Error(t, err) + require.Contains(t, err.Error(), "create otlp trace exporter") +} + +// TestShutdown_nilProvider confirms Shutdown is safe on a zero-value OTLP. +func TestShutdown_nilProvider(t *testing.T) { + var o OTLP + require.NoError(t, o.Shutdown(context.Background())) +} diff --git a/output/adapter.go b/output/adapter.go index ae26e1c..b769d6b 100644 --- a/output/adapter.go +++ b/output/adapter.go @@ -4,8 +4,13 @@ import ( "context" "github.com/observiq/blitz/embed" + "go.opentelemetry.io/otel/trace" ) +// adapterTracerScope is the instrumentation scope for the per-emit-cycle spans +// the consumer adapters create when TelemetrySettings.PerBatchSpans is on. +const adapterTracerScope = "github.com/observiq/blitz/output" + // WriterAsLogConsumer wraps a Writer so it can be used in contexts that // expect an embed.LogConsumer. The adapter pushes each record in the // batch through Writer.Write in order, returning the first error it @@ -13,23 +18,30 @@ import ( // // CLI generator wiring uses this adapter to bridge migrated modules // (which talk to embed.LogConsumer) with the existing Output instances -// (which implement Writer). +// (which implement Writer). tel carries blitz's self-telemetry providers; +// when tel.PerBatchSpans is set, each ConsumeLogs call is wrapped in a span. // // Panics on nil writer — a nil writer is a programming bug, not a // runtime condition, and catching it at construction surfaces the // failure at the boundary rather than deep in ConsumeLogs. -func WriterAsLogConsumer(w Writer) embed.LogConsumer { +func WriterAsLogConsumer(w Writer, tel embed.TelemetrySettings) embed.LogConsumer { if w == nil { panic("output.WriterAsLogConsumer: writer cannot be nil") } - return &writerAsLogConsumer{w: w} + return &writerAsLogConsumer{w: w, tel: tel} } type writerAsLogConsumer struct { - w Writer + w Writer + tel embed.TelemetrySettings } func (a *writerAsLogConsumer) ConsumeLogs(ctx context.Context, records []embed.LogRecord) error { + if a.tel.PerBatchSpans { + var span trace.Span + ctx, span = a.tel.Tracer(adapterTracerScope).Start(ctx, "blitz.emit.logs") + defer span.End() + } for i := range records { if err := a.w.Write(ctx, records[i]); err != nil { return err @@ -45,21 +57,27 @@ func (a *writerAsLogConsumer) ConsumeLogs(ctx context.Context, records []embed.L // // Standalone CLI metric-generator wiring uses this adapter to bridge // modules that talk to embed.MetricConsumer with existing outputs that -// implement MetricWriter. +// implement MetricWriter. tel.PerBatchSpans wraps each call in a span. // // Panics on nil writer; see WriterAsLogConsumer for the rationale. -func WriterAsMetricConsumer(w MetricWriter) embed.MetricConsumer { +func WriterAsMetricConsumer(w MetricWriter, tel embed.TelemetrySettings) embed.MetricConsumer { if w == nil { panic("output.WriterAsMetricConsumer: writer cannot be nil") } - return &writerAsMetricConsumer{w: w} + return &writerAsMetricConsumer{w: w, tel: tel} } type writerAsMetricConsumer struct { - w MetricWriter + w MetricWriter + tel embed.TelemetrySettings } func (a *writerAsMetricConsumer) ConsumeMetrics(ctx context.Context, points []embed.MetricPoint) error { + if a.tel.PerBatchSpans { + var span trace.Span + ctx, span = a.tel.Tracer(adapterTracerScope).Start(ctx, "blitz.emit.metrics") + defer span.End() + } for i := range points { if err := a.w.WriteMetric(ctx, points[i]); err != nil { return err @@ -75,23 +93,29 @@ func (a *writerAsMetricConsumer) ConsumeMetrics(ctx context.Context, points []em // // Standalone CLI trace-generator wiring uses this adapter to bridge // modules that talk to embed.TraceConsumer with existing outputs that -// implement TraceWriter. +// implement TraceWriter. tel.PerBatchSpans wraps each call in a span. // // Panics on nil writer — a nil writer is a programming bug, not a // runtime condition, and catching it at construction surfaces the // failure at the boundary rather than deep in ConsumeTraces. -func WriterAsTraceConsumer(w TraceWriter) embed.TraceConsumer { +func WriterAsTraceConsumer(w TraceWriter, tel embed.TelemetrySettings) embed.TraceConsumer { if w == nil { panic("output.WriterAsTraceConsumer: writer cannot be nil") } - return &writerAsTraceConsumer{w: w} + return &writerAsTraceConsumer{w: w, tel: tel} } type writerAsTraceConsumer struct { - w TraceWriter + w TraceWriter + tel embed.TelemetrySettings } func (a *writerAsTraceConsumer) ConsumeTraces(ctx context.Context, spans []embed.Span) error { + if a.tel.PerBatchSpans { + var span trace.Span + ctx, span = a.tel.Tracer(adapterTracerScope).Start(ctx, "blitz.emit.traces") + defer span.End() + } for i := range spans { if err := a.w.WriteTrace(ctx, spans[i]); err != nil { return err diff --git a/output/adapter_test.go b/output/adapter_test.go index 61de34a..37699b3 100644 --- a/output/adapter_test.go +++ b/output/adapter_test.go @@ -8,6 +8,8 @@ import ( "github.com/observiq/blitz/embed" "github.com/observiq/blitz/output" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" ) type recordingWriter struct { @@ -28,7 +30,7 @@ func (w *recordingWriter) Write(_ context.Context, rec output.LogRecord) error { func TestWriterAsLogConsumerPushesEachRecord(t *testing.T) { w := &recordingWriter{} - c := output.WriterAsLogConsumer(w) + c := output.WriterAsLogConsumer(w, embed.NopTelemetry()) batch := []embed.LogRecord{ {Message: "one"}, @@ -48,10 +50,54 @@ func TestWriterAsLogConsumerPushesEachRecord(t *testing.T) { } } +func TestWriterAsLogConsumerPerBatchSpanWhenEnabled(t *testing.T) { + exporter := tracetest.NewInMemoryExporter() + tp := sdktrace.NewTracerProvider(sdktrace.WithSyncer(exporter)) + + w := &recordingWriter{} + c := output.WriterAsLogConsumer(w, embed.TelemetrySettings{TracerProvider: tp, PerBatchSpans: true}) + + batch := []embed.LogRecord{ + {Message: "one"}, + {Message: "two"}, + } + if err := c.ConsumeLogs(context.Background(), batch); err != nil { + t.Fatalf("ConsumeLogs: %v", err) + } + + spans := exporter.GetSpans() + if got, want := len(spans), 1; got != want { + t.Fatalf("exported %d spans, want %d", got, want) + } + if got, want := spans[0].Name, "blitz.emit.logs"; got != want { + t.Errorf("span name %q, want %q", got, want) + } +} + +func TestWriterAsLogConsumerNoSpanWhenDisabled(t *testing.T) { + exporter := tracetest.NewInMemoryExporter() + tp := sdktrace.NewTracerProvider(sdktrace.WithSyncer(exporter)) + + w := &recordingWriter{} + c := output.WriterAsLogConsumer(w, embed.TelemetrySettings{TracerProvider: tp}) + + batch := []embed.LogRecord{ + {Message: "one"}, + {Message: "two"}, + } + if err := c.ConsumeLogs(context.Background(), batch); err != nil { + t.Fatalf("ConsumeLogs: %v", err) + } + + if got := exporter.GetSpans(); len(got) != 0 { + t.Fatalf("exported %d spans, want 0", len(got)) + } +} + func TestWriterAsLogConsumerStopsOnFirstError(t *testing.T) { wantErr := errors.New("boom") w := &recordingWriter{err: wantErr} - c := output.WriterAsLogConsumer(w) + c := output.WriterAsLogConsumer(w, embed.NopTelemetry()) err := c.ConsumeLogs(context.Background(), []embed.LogRecord{ {Message: "one"}, @@ -80,7 +126,7 @@ func (w *recordingMetricWriter) WriteMetric(_ context.Context, rec output.Metric func TestWriterAsMetricConsumerPushesEachPoint(t *testing.T) { w := &recordingMetricWriter{} - c := output.WriterAsMetricConsumer(w) + c := output.WriterAsMetricConsumer(w, embed.NopTelemetry()) batch := []embed.MetricPoint{ {Name: "one"}, @@ -103,7 +149,7 @@ func TestWriterAsMetricConsumerPushesEachPoint(t *testing.T) { func TestWriterAsMetricConsumerStopsOnFirstError(t *testing.T) { wantErr := errors.New("metric boom") w := &recordingMetricWriter{err: wantErr} - c := output.WriterAsMetricConsumer(w) + c := output.WriterAsMetricConsumer(w, embed.NopTelemetry()) err := c.ConsumeMetrics(context.Background(), []embed.MetricPoint{ {Name: "one"}, @@ -132,7 +178,7 @@ func (w *recordingTraceWriter) WriteTrace(_ context.Context, rec output.TraceRec func TestWriterAsTraceConsumerPushesEachSpan(t *testing.T) { w := &recordingTraceWriter{} - c := output.WriterAsTraceConsumer(w) + c := output.WriterAsTraceConsumer(w, embed.NopTelemetry()) batch := []embed.Span{ {Name: "one"}, @@ -155,7 +201,7 @@ func TestWriterAsTraceConsumerPushesEachSpan(t *testing.T) { func TestWriterAsTraceConsumerStopsOnFirstError(t *testing.T) { wantErr := errors.New("trace boom") w := &recordingTraceWriter{err: wantErr} - c := output.WriterAsTraceConsumer(w) + c := output.WriterAsTraceConsumer(w, embed.NopTelemetry()) err := c.ConsumeTraces(context.Background(), []embed.Span{ {Name: "one"}, @@ -172,7 +218,7 @@ func TestWriterAsTraceConsumerPanicsOnNilWriter(t *testing.T) { t.Fatal("expected panic on nil writer, got none") } }() - _ = output.WriterAsTraceConsumer(nil) + _ = output.WriterAsTraceConsumer(nil, embed.NopTelemetry()) } func TestWriterAsLogConsumerPanicsOnNilWriter(t *testing.T) { @@ -181,7 +227,7 @@ func TestWriterAsLogConsumerPanicsOnNilWriter(t *testing.T) { t.Fatal("expected panic on nil writer, got none") } }() - _ = output.WriterAsLogConsumer(nil) + _ = output.WriterAsLogConsumer(nil, embed.NopTelemetry()) } func TestWriterAsMetricConsumerPanicsOnNilWriter(t *testing.T) { @@ -190,5 +236,5 @@ func TestWriterAsMetricConsumerPanicsOnNilWriter(t *testing.T) { t.Fatal("expected panic on nil writer, got none") } }() - _ = output.WriterAsMetricConsumer(nil) + _ = output.WriterAsMetricConsumer(nil, embed.NopTelemetry()) } diff --git a/package/completions/blitz.bash b/package/completions/blitz.bash index bdcf4c0..3e59c0d 100644 --- a/package/completions/blitz.bash +++ b/package/completions/blitz.bash @@ -632,6 +632,10 @@ _blitz_help() two_word_flags+=("--output-udp-port") flags+=("--output-udp-workers=") two_word_flags+=("--output-udp-workers") + flags+=("--telemetry-traces-insecure") + flags+=("--telemetry-traces-otlpendpoint=") + two_word_flags+=("--telemetry-traces-otlpendpoint") + flags+=("--telemetry-traces-perbatchspans") must_have_one_flag=() must_have_one_noun=() @@ -911,6 +915,10 @@ _blitz_version() two_word_flags+=("--output-udp-port") flags+=("--output-udp-workers=") two_word_flags+=("--output-udp-workers") + flags+=("--telemetry-traces-insecure") + flags+=("--telemetry-traces-otlpendpoint=") + two_word_flags+=("--telemetry-traces-otlpendpoint") + flags+=("--telemetry-traces-perbatchspans") must_have_one_flag=() must_have_one_noun=() @@ -1191,6 +1199,10 @@ _blitz_root_command() two_word_flags+=("--output-udp-port") flags+=("--output-udp-workers=") two_word_flags+=("--output-udp-workers") + flags+=("--telemetry-traces-insecure") + flags+=("--telemetry-traces-otlpendpoint=") + two_word_flags+=("--telemetry-traces-otlpendpoint") + flags+=("--telemetry-traces-perbatchspans") must_have_one_flag=() must_have_one_noun=()