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
17 changes: 16 additions & 1 deletion generator/fix/fix.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -96,6 +99,7 @@ type Generator struct {
logger *zap.Logger
cfg Config
consumer embed.LogConsumer
metrics *generator.Metrics

wg sync.WaitGroup
stopCh chan struct{}
Expand All @@ -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")
}
Expand All @@ -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
}
Expand All @@ -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)
Expand All @@ -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{})
Expand Down Expand Up @@ -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()
}
}
Expand Down
127 changes: 116 additions & 11 deletions generator/fix/fix_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand All @@ -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()))
Expand All @@ -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)
Expand Down Expand Up @@ -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()))
Expand All @@ -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()))
Expand Down Expand Up @@ -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()))
Expand All @@ -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
}
6 changes: 3 additions & 3 deletions internal/dispatch/embed.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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)
}
Loading