Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion cmd/blitz/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -190,14 +190,15 @@ 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
}
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))
Expand Down
33 changes: 25 additions & 8 deletions output/nop/nop.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
78 changes: 78 additions & 0 deletions output/nop/nop_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
13 changes: 13 additions & 0 deletions output/stdout/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@ package stdout
import (
"fmt"
"time"

"github.com/observiq/blitz/embed"
)

const defaultFlushInterval = 100 * time.Millisecond
Expand All @@ -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.
Expand Down
21 changes: 20 additions & 1 deletion output/stdout/stdout.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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{}),
Expand All @@ -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()

Expand Down Expand Up @@ -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
}
Expand Down
72 changes: 72 additions & 0 deletions output/stdout/stdout_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down Expand Up @@ -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")
}
Loading