From f6d4081bf2783c9f8f7363c61e3cefdd3caa0a9e Mon Sep 17 00:00:00 2001 From: Dylan Myers Date: Thu, 6 Aug 2026 09:20:56 -0400 Subject: [PATCH] feat(o11y): wire fix generator into the self-telemetry spine (PIPE-1066) Assisted-by: Claude Opus 4.8 --- generator/fix/fix.go | 17 ++++- generator/fix/fix_test.go | 127 +++++++++++++++++++++++++++++++++---- internal/dispatch/embed.go | 6 +- 3 files changed, 135 insertions(+), 15 deletions(-) diff --git a/generator/fix/fix.go b/generator/fix/fix.go index cb9204f..00e3eb7 100644 --- a/generator/fix/fix.go +++ b/generator/fix/fix.go @@ -29,9 +29,12 @@ import ( "sync" "time" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/metric" "go.uber.org/zap" "github.com/observiq/blitz/embed" + "github.com/observiq/blitz/generator" "github.com/observiq/blitz/generator/fix/catalog" "github.com/observiq/blitz/generator/fix/catalog/v44/app" "github.com/observiq/blitz/generator/fix/state" @@ -96,6 +99,7 @@ type Generator struct { logger *zap.Logger cfg Config consumer embed.LogConsumer + metrics *generator.Metrics wg sync.WaitGroup stopCh chan struct{} @@ -105,7 +109,7 @@ type Generator struct { // FIX message as a size-1 batch via ConsumeLogs. Returns an error for // invalid inputs (nil logger, nil consumer, workers < 1, non-positive // rate). -func New(logger *zap.Logger, cfg Config, consumer embed.LogConsumer) (*Generator, error) { +func New(logger *zap.Logger, cfg Config, consumer embed.LogConsumer, tel embed.TelemetrySettings) (*Generator, error) { if logger == nil { return nil, fmt.Errorf("logger cannot be nil") } @@ -130,10 +134,15 @@ func New(logger *zap.Logger, cfg Config, consumer embed.LogConsumer) (*Generator if len(cfg.EnabledCategories) == 0 { cfg.EnabledCategories = catalog.AllAssetCategories() } + metrics, err := generator.NewMetrics(tel.MeterProvider) + if err != nil { + return nil, fmt.Errorf("build generator metrics: %w", err) + } return &Generator{ logger: logger, cfg: cfg, consumer: consumer, + metrics: metrics, stopCh: make(chan struct{}), }, nil } @@ -148,6 +157,7 @@ func (g *Generator) Start(_ context.Context) error { zap.Duration("rate", g.cfg.Rate), zap.String("version", g.cfg.Version.String()), ) + g.metrics.BlitzGeneratorActiveWorkersGauge.Record(context.Background(), int64(g.cfg.Workers), componentName) for i := 0; i < g.cfg.Workers; i++ { g.wg.Add(1) go g.runWorker(i) @@ -159,6 +169,7 @@ func (g *Generator) Start(_ context.Context) error { // when all workers have stopped or ctx is canceled. func (g *Generator) Stop(ctx context.Context) error { g.logger.Info("Stopping FIX generator") + g.metrics.BlitzGeneratorActiveWorkersGauge.Record(context.Background(), 0, componentName) close(g.stopCh) done := make(chan struct{}) @@ -214,7 +225,11 @@ func (g *Generator) runWorker(workerIdx int) { } if err := g.consumer.ConsumeLogs(ctx, []embed.LogRecord{rec}); err != nil { g.logger.Debug("FIX emit failed", zap.Error(err)) + g.metrics.BlitzGeneratorWriteErrorsCounter.Add(context.Background(), 1, componentName, + metric.WithAttributeSet(attribute.NewSet(attribute.String("error_type", "consume"))), + ) } + g.metrics.BlitzGeneratorEntriesCounter.Add(context.Background(), 1, componentName) cancel() } } diff --git a/generator/fix/fix_test.go b/generator/fix/fix_test.go index a89d6ae..1400b0d 100644 --- a/generator/fix/fix_test.go +++ b/generator/fix/fix_test.go @@ -3,12 +3,18 @@ package fix import ( "bytes" "context" + "errors" "sync" "testing" "time" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "go.opentelemetry.io/otel/log/logtest" + "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" "github.com/observiq/blitz/embed" @@ -57,6 +63,105 @@ func (c *captureConsumer) Count() int { return len(c.got) } +// failingConsumer always returns an error, to exercise the write-error path. +type failingConsumer struct { + mu sync.Mutex + calls int +} + +func (c *failingConsumer) ConsumeLogs(context.Context, []embed.LogRecord) error { + c.mu.Lock() + c.calls++ + c.mu.Unlock() + return errors.New("consume failed") +} + +func (c *failingConsumer) Count() int { + c.mu.Lock() + defer c.mu.Unlock() + return c.calls +} + +func hasMetric(rm metricdata.ResourceMetrics, name string) bool { + for _, sm := range rm.ScopeMetrics { + for _, m := range sm.Metrics { + if m.Name == name { + return true + } + } + } + return false +} + +// TestFIX_SelfTelemetry confirms the FIX generator routes its own metrics +// through the injected MeterProvider and bridges its logs into the injected +// LoggerProvider, matching every other embed-eligible generator. The failing +// consumer also exercises the write-error metric path. +func TestFIX_SelfTelemetry(t *testing.T) { + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + logRec := logtest.NewRecorder() + tel := embed.TelemetrySettings{ + Logger: zap.NewNop(), + MeterProvider: mp, + LoggerProvider: logRec, + } + + cons := &failingConsumer{} + // A caller constructing a generator directly bridges the logger itself, + // as config.LoadModules does; blitz no longer re-bridges per component. + g, err := New(tel.BridgedLogger(zap.NewNop()), Config{Workers: 1, Rate: 5 * time.Millisecond}, cons, tel) + require.NoError(t, err) + require.NoError(t, g.Start(context.Background())) + + require.Eventually(t, + func() bool { return cons.Count() >= 1 }, + 2*time.Second, 10*time.Millisecond, + "expected the generator to attempt at least one emit", + ) + + stopCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + require.NoError(t, g.Stop(stopCtx)) + + // Self-metrics reached the injected provider. + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &rm)) + require.True(t, hasMetric(rm, "blitz.generator.active_workers"), "active_workers gauge") + require.True(t, hasMetric(rm, "blitz.generator.entries"), "entries counter") + require.True(t, hasMetric(rm, "blitz.generator.write_errors"), "write_errors counter") + + // Internal logs bridged into the injected LoggerProvider. + require.NotEmpty(t, logRec.Result(), "expected bridged log records") +} + +// failingMeter embeds a no-op Meter and overrides one instrument constructor to +// return an error, so generator.NewMetrics fails. +type failingMeter struct { + metric.Meter +} + +func (failingMeter) Int64Gauge(string, ...metric.Int64GaugeOption) (metric.Int64Gauge, error) { + return nil, errors.New("instrument error") +} + +// failingMeterProvider hands out a failingMeter. +type failingMeterProvider struct { + metric.MeterProvider +} + +func (failingMeterProvider) Meter(string, ...metric.MeterOption) metric.Meter { + return failingMeter{Meter: metricnoop.NewMeterProvider().Meter("test")} +} + +// TestNew_MetricsBuildError covers the metric-construction error path in New. +func TestNew_MetricsBuildError(t *testing.T) { + _, err := New(zap.NewNop(), DefaultConfig(), &captureConsumer{}, + embed.TelemetrySettings{MeterProvider: failingMeterProvider{}}) + require.Error(t, err) + require.Contains(t, err.Error(), "build generator metrics") +} + func TestDefaultConfig(t *testing.T) { c := DefaultConfig() assert.Equal(t, 1, c.Workers) @@ -65,32 +170,32 @@ func TestDefaultConfig(t *testing.T) { } func TestNewRejectsNilLogger(t *testing.T) { - _, err := New(nil, DefaultConfig(), &captureConsumer{}) + _, err := New(nil, DefaultConfig(), &captureConsumer{}, embed.NopTelemetry()) require.Error(t, err) } func TestNewRejectsNilConsumer(t *testing.T) { - _, err := New(zap.NewNop(), DefaultConfig(), nil) + _, err := New(zap.NewNop(), DefaultConfig(), nil, embed.NopTelemetry()) require.Error(t, err) } func TestNewRejectsZeroWorkers(t *testing.T) { cfg := DefaultConfig() cfg.Workers = 0 - _, err := New(zap.NewNop(), cfg, &captureConsumer{}) + _, err := New(zap.NewNop(), cfg, &captureConsumer{}, embed.NopTelemetry()) require.Error(t, err) } func TestNewRejectsNonPositiveRate(t *testing.T) { cfg := DefaultConfig() cfg.Rate = 0 - _, err := New(zap.NewNop(), cfg, &captureConsumer{}) + _, err := New(zap.NewNop(), cfg, &captureConsumer{}, embed.NopTelemetry()) require.Error(t, err) } func TestNewDefaultsVersionAndCompIDs(t *testing.T) { cfg := Config{Workers: 1, Rate: time.Second} - g, err := New(zap.NewNop(), cfg, &captureConsumer{}) + g, err := New(zap.NewNop(), cfg, &captureConsumer{}, embed.NopTelemetry()) require.NoError(t, err) assert.Equal(t, catalog.V44, g.cfg.Version) assert.Equal(t, "BLITZ", g.cfg.SenderCompID) @@ -105,7 +210,7 @@ func TestEmitsMessagesAtRate(t *testing.T) { Rate: 20 * time.Millisecond, Version: catalog.V44, Seed: 42, - }, cons) + }, cons, embed.NopTelemetry()) require.NoError(t, err) require.NoError(t, g.Start(context.Background())) @@ -129,7 +234,7 @@ func TestGoldenOutputDeterministicFromSeed(t *testing.T) { runOnce := func() [][]byte { cons := &captureConsumer{} - g, err := New(zap.NewNop(), cfg, cons) + g, err := New(zap.NewNop(), cfg, cons, embed.NopTelemetry()) require.NoError(t, err) require.NoError(t, g.Start(context.Background())) require.Eventually(t, func() bool { return cons.Count() >= 5 }, 2*time.Second, 5*time.Millisecond) @@ -160,7 +265,7 @@ func TestV50SP2EmitsFIXTBeginString(t *testing.T) { Rate: 20 * time.Millisecond, Version: catalog.V50SP2, Seed: 42, - }, cons) + }, cons, embed.NopTelemetry()) require.NoError(t, err) require.NoError(t, g.Start(context.Background())) @@ -182,7 +287,7 @@ func TestV42EmitsFIX42BeginString(t *testing.T) { Rate: 20 * time.Millisecond, Version: catalog.V42, Seed: 42, - }, cons) + }, cons, embed.NopTelemetry()) require.NoError(t, err) require.NoError(t, g.Start(context.Background())) @@ -219,7 +324,7 @@ func TestEmitsResourceWithHostNameAndFixVersion(t *testing.T) { Rate: 20 * time.Millisecond, Version: tc.version, Seed: 42, - }, cons) + }, cons, embed.NopTelemetry()) require.NoError(t, err) require.NoError(t, g.Start(context.Background())) @@ -241,7 +346,7 @@ func TestEmitsResourceWithHostNameAndFixVersion(t *testing.T) { // *Generator is embed-eligible. If embed.ProducerMarker is removed or // the Module interface changes, this fails to compile. func TestGeneratorSatisfiesProducerModule(t *testing.T) { - g, err := New(zap.NewNop(), DefaultConfig(), &captureConsumer{}) + g, err := New(zap.NewNop(), DefaultConfig(), &captureConsumer{}, embed.NopTelemetry()) require.NoError(t, err) var _ embed.ProducerModule = g } diff --git a/internal/dispatch/embed.go b/internal/dispatch/embed.go index 81a297d..334edb5 100644 --- a/internal/dispatch/embed.go +++ b/internal/dispatch/embed.go @@ -153,7 +153,7 @@ func ForEmbed(logger *zap.Logger, genCfg config.Generator, consumers EmbedConsum if err := consumers.requireLog(genCfg.Type); err != nil { return nil, err } - return newFIX(logger, genCfg.FIX, consumers.LogConsumer) + return newFIX(logger, genCfg.FIX, consumers.LogConsumer, tel) case config.GeneratorTypeHostMetrics: if err := consumers.requireMetric(genCfg.Type); err != nil { return nil, err @@ -214,7 +214,7 @@ func yamlSeedDefault(yamlSeed int64) int64 { // catalog-typed fix.Config and constructs a FIX generator. Version and // EnabledCategories strings are validated; an empty version defaults to // FIX 4.4 and an empty EnabledCategories means "all 10 categories". -func newFIX(logger *zap.Logger, cfg config.FIXGeneratorConfig, consumer embed.LogConsumer) (embed.ProducerModule, error) { +func newFIX(logger *zap.Logger, cfg config.FIXGeneratorConfig, consumer embed.LogConsumer, tel embed.TelemetrySettings) (embed.ProducerModule, error) { fc := fixgen.Config{ Workers: cfg.Workers, Rate: cfg.Rate, @@ -236,5 +236,5 @@ func newFIX(logger *zap.Logger, cfg config.FIXGeneratorConfig, consumer embed.Lo } fc.EnabledCategories = append(fc.EnabledCategories, c) } - return fixgen.New(logger, fc, consumer) + return fixgen.New(logger, fc, consumer, tel) }