diff --git a/cmd/blitz/main.go b/cmd/blitz/main.go index d67d696..f719fbe 100644 --- a/cmd/blitz/main.go +++ b/cmd/blitz/main.go @@ -190,7 +190,7 @@ func run(cmd *cobra.Command, args []string) error { var outputInstance output.Output switch cfg.Output.Type { case config.OutputTypeNop: - outputInstance, err = nop.New(logger) + outputInstance, err = nop.New(logger, tel) if err != nil { logger.Error("Failed to create NOP output", zap.Error(err)) return err @@ -198,6 +198,7 @@ func run(cmd *cobra.Command, args []string) error { case config.OutputTypeStdout: outputInstance, err = stdoutout.New(logger, stdoutout.WithFlushInterval(cfg.Output.Stdout.FlushInterval), + stdoutout.WithTelemetry(tel), ) if err != nil { logger.Error("Failed to create stdout output", zap.Error(err)) diff --git a/output/nop/nop.go b/output/nop/nop.go index 339b4ee..09e37f6 100644 --- a/output/nop/nop.go +++ b/output/nop/nop.go @@ -4,35 +4,52 @@ import ( "context" "fmt" + "github.com/observiq/blitz/embed" "github.com/observiq/blitz/output" "github.com/observiq/blitz/telemetry" "go.uber.org/zap" ) -// NopOutput is a no-operation output that performs no work +// outputType is the output_type attribute value for nop metrics. +const outputType = "nop" + +// NopOutput is a no-operation output that discards every record. It still +// records how many records it received, so a load run can measure how much was +// pushed into the void, and it bridges its logger like every other output. type NopOutput struct { - logger *zap.Logger + logger *zap.Logger + metrics *output.Metrics } -// New creates a new no-operation output -func New(logger *zap.Logger) (*NopOutput, error) { +// New creates a new no-operation output. tel carries blitz's self-telemetry +// providers: metrics route through tel.MeterProvider and the logger is bridged +// into tel.LoggerProvider. +func New(logger *zap.Logger, tel embed.TelemetrySettings) (*NopOutput, error) { if logger == nil { return nil, fmt.Errorf("logger cannot be nil") } + m, err := output.NewMetrics(tel.MeterProvider) + if err != nil { + return nil, fmt.Errorf("build output metrics: %w", err) + } + return &NopOutput{ - logger: logger.Named("output-nop"), + logger: logger.Named("output-nop"), + metrics: m, }, nil } -// Write performs no work (data is discarded) +// Write discards the record after counting it. The write is a synchronous +// no-op, already bracketed by the consumer adapter's emit span, so it gets no +// span of its own. func (o *NopOutput) Write(ctx context.Context, data output.LogRecord) error { - // No-op: data is discarded + o.metrics.BlitzOutputEntriesReceivedCounter.Add(ctx, 1, outputType, "logs") return nil } // Stop performs no work -func (o *NopOutput) Stop(ctx context.Context) error { +func (o *NopOutput) Stop(_ context.Context) error { o.logger.Info("Stopping NOP output") return nil } diff --git a/output/nop/nop_test.go b/output/nop/nop_test.go new file mode 100644 index 0000000..3202a5a --- /dev/null +++ b/output/nop/nop_test.go @@ -0,0 +1,78 @@ +package nop + +import ( + "context" + "errors" + "testing" + + "github.com/observiq/blitz/embed" + "github.com/observiq/blitz/output" + "github.com/observiq/blitz/telemetry" + "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/metric" + metricnoop "go.opentelemetry.io/otel/metric/noop" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" + "go.uber.org/zap" +) + +func TestNew_NilLogger(t *testing.T) { + _, err := New(nil, embed.NopTelemetry()) + require.Error(t, err) +} + +func TestSupportedTelemetry(t *testing.T) { + o, err := New(zap.NewNop(), embed.NopTelemetry()) + require.NoError(t, err) + require.Equal(t, []telemetry.Type{telemetry.Logs}, o.SupportedTelemetry()) +} + +// failingMeter overrides the first instrument the output registry builds +// (an Int64Gauge) to return an error, so output.NewMetrics fails. +type failingMeter struct { + metric.Meter +} + +func (failingMeter) Int64Gauge(string, ...metric.Int64GaugeOption) (metric.Int64Gauge, error) { + return nil, errors.New("instrument error") +} + +type failingMeterProvider struct { + metric.MeterProvider +} + +func (failingMeterProvider) Meter(string, ...metric.MeterOption) metric.Meter { + return failingMeter{Meter: metricnoop.NewMeterProvider().Meter("test")} +} + +func TestNew_MetricsError(t *testing.T) { + _, err := New(zap.NewNop(), embed.TelemetrySettings{MeterProvider: failingMeterProvider{}}) + require.Error(t, err) + require.Contains(t, err.Error(), "build output metrics") +} + +// TestWrite_countsAndDiscards confirms nop counts every record through the +// injected MeterProvider while discarding the data. +func TestWrite_countsAndDiscards(t *testing.T) { + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + + o, err := New(zap.NewNop(), embed.TelemetrySettings{MeterProvider: mp}) + require.NoError(t, err) + + require.NoError(t, o.Write(context.Background(), output.LogRecord{Message: "discarded"})) + require.NoError(t, o.Stop(context.Background())) + + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &rm)) + + found := false + for _, sm := range rm.ScopeMetrics { + for _, m := range sm.Metrics { + if m.Name == "blitz.output.entries_received" { + found = true + } + } + } + require.True(t, found, "expected blitz.output.entries_received counter") +} diff --git a/output/stdout/options.go b/output/stdout/options.go index 6e28d7d..2a16121 100644 --- a/output/stdout/options.go +++ b/output/stdout/options.go @@ -3,6 +3,8 @@ package stdout import ( "fmt" "time" + + "github.com/observiq/blitz/embed" ) const defaultFlushInterval = 100 * time.Millisecond @@ -12,6 +14,17 @@ type Option func(*config) error type config struct { flushInterval time.Duration + tel embed.TelemetrySettings +} + +// WithTelemetry sets the OTel providers blitz routes its self-telemetry +// through: metrics via tel.MeterProvider, the log bridge via tel.LoggerProvider, +// and the gated flush span via tel.TracerProvider. +func WithTelemetry(tel embed.TelemetrySettings) Option { + return func(c *config) error { + c.tel = tel + return nil + } } // WithFlushInterval sets the interval at which the internal buffer is flushed to stdout. diff --git a/output/stdout/stdout.go b/output/stdout/stdout.go index 0dfd873..cf14880 100644 --- a/output/stdout/stdout.go +++ b/output/stdout/stdout.go @@ -9,16 +9,22 @@ import ( "sync" "time" + "github.com/observiq/blitz/embed" "github.com/observiq/blitz/output" "github.com/observiq/blitz/telemetry" "go.uber.org/zap" ) +// outputType is the output_type attribute value for stdout metrics. +const outputType = "stdout" + // StdoutOutput writes log records to standard output via a buffered writer. // Records are batched in memory and flushed to os.Stdout periodically, reducing // per-record syscall overhead under high worker counts. type StdoutOutput struct { logger *zap.Logger + tel embed.TelemetrySettings + metrics *output.Metrics writer *bufio.Writer mu sync.Mutex flushInterval time.Duration @@ -43,8 +49,15 @@ func New(logger *zap.Logger, opts ...Option) (*StdoutOutput, error) { } } + m, err := output.NewMetrics(cfg.tel.MeterProvider) + if err != nil { + return nil, fmt.Errorf("build output metrics: %w", err) + } + o := &StdoutOutput{ logger: logger.Named("output-stdout"), + tel: cfg.tel, + metrics: m, writer: bufio.NewWriter(os.Stdout), flushInterval: cfg.flushInterval, stopCh: make(chan struct{}), @@ -57,7 +70,9 @@ func New(logger *zap.Logger, opts ...Option) (*StdoutOutput, error) { } // Write buffers the log record for the next flush. -func (o *StdoutOutput) Write(_ context.Context, data output.LogRecord) error { +func (o *StdoutOutput) Write(ctx context.Context, data output.LogRecord) error { + o.metrics.BlitzOutputEntriesReceivedCounter.Add(ctx, 1, outputType, "logs") + o.mu.Lock() defer o.mu.Unlock() @@ -93,9 +108,13 @@ func (o *StdoutOutput) flushLoop() { for { select { case <-ticker.C: + // The flush is the actual stdout I/O, batching many buffered + // records, so it gets a standalone gated span. + _, span := output.StartSendSpan(context.Background(), o.tel, "blitz.output.stdout.flush") o.mu.Lock() _ = o.writer.Flush() o.mu.Unlock() + span.End() case <-o.stopCh: return } diff --git a/output/stdout/stdout_test.go b/output/stdout/stdout_test.go index cc902e1..13de22b 100644 --- a/output/stdout/stdout_test.go +++ b/output/stdout/stdout_test.go @@ -10,10 +10,18 @@ import ( "testing" "time" + "github.com/observiq/blitz/embed" "github.com/observiq/blitz/output" "github.com/observiq/blitz/telemetry" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/metric" + metricnoop "go.opentelemetry.io/otel/metric/noop" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" + sdktrace "go.opentelemetry.io/otel/sdk/trace" + "go.opentelemetry.io/otel/sdk/trace/tracetest" + "go.uber.org/zap" "go.uber.org/zap/zaptest" ) @@ -281,3 +289,67 @@ func TestStdoutOutput_Stop(t *testing.T) { err = out.Stop(context.Background()) assert.NoError(t, err) } + +// failingMeter overrides the first instrument the output registry builds (an +// Int64Gauge) to error, so output.NewMetrics fails. +type failingMeter struct{ metric.Meter } + +func (failingMeter) Int64Gauge(string, ...metric.Int64GaugeOption) (metric.Int64Gauge, error) { + return nil, fmt.Errorf("instrument error") +} + +type failingMeterProvider struct{ metric.MeterProvider } + +func (failingMeterProvider) Meter(string, ...metric.MeterOption) metric.Meter { + return failingMeter{Meter: metricnoop.NewMeterProvider().Meter("test")} +} + +func TestStdout_MetricsError(t *testing.T) { + _, err := New(zap.NewNop(), WithTelemetry(embed.TelemetrySettings{MeterProvider: failingMeterProvider{}})) + require.Error(t, err) + require.Contains(t, err.Error(), "build output metrics") +} + +// TestStdout_SelfTelemetry confirms stdout records its entries metric through +// the injected MeterProvider and emits a gated flush span through the injected +// TracerProvider. +func TestStdout_SelfTelemetry(t *testing.T) { + oldStdout := os.Stdout + _, w, err := os.Pipe() + require.NoError(t, err) + os.Stdout = w + defer func() { os.Stdout = oldStdout; _ = w.Close() }() + + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + spanRec := tracetest.NewSpanRecorder() + tp := sdktrace.NewTracerProvider(sdktrace.WithSpanProcessor(spanRec)) + tel := embed.TelemetrySettings{MeterProvider: mp, TracerProvider: tp, PerBatchSpans: true} + + out, err := New(zap.NewNop(), WithFlushInterval(10*time.Millisecond), WithTelemetry(tel)) + require.NoError(t, err) + require.NoError(t, out.Write(context.Background(), output.LogRecord{Message: "x"})) + + require.Eventually(t, func() bool { + for _, s := range spanRec.Ended() { + if s.Name() == "blitz.output.stdout.flush" { + return true + } + } + return false + }, 2*time.Second, 10*time.Millisecond, "expected a flush span") + + require.NoError(t, out.Stop(context.Background())) + + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &rm)) + found := false + for _, sm := range rm.ScopeMetrics { + for _, m := range sm.Metrics { + if m.Name == "blitz.output.entries_received" { + found = true + } + } + } + require.True(t, found, "expected blitz.output.entries_received counter") +}