From 1f311adf769b926d472d99c741535de8757a3715 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adnan=20Rahi=C4=87?= Date: Wed, 7 Oct 2026 13:52:18 +0200 Subject: [PATCH 1/2] fix(otlp-grpc): send generator resource and real severity on logs, metrics and traces The OTLP gRPC output dropped the resource every generator attaches (Metadata.Resource): logs and traces were always sent with a hard-coded service.name=blitz resource, and metrics with buildMetricRequest(metrics, nil). host.name, telemetry.source, apache.format and the rest never reached the collector, so pipelines could not route or filter on them. Per-record Metadata.Attributes on logs were dropped too. Each queued log record, metric and span now carries its generator's resource. A batch is split into one ResourceLogs / ResourceMetrics / ResourceSpans per distinct resource, in order of first appearance. service.name=blitz is still added unless the generator sets its own. mapSeverityNumber only matched exact uppercase DEBUG/INFO/WARN/ERROR/FATAL, so lowercase levels (kubernetes, apache error), Windows Event Log level names (Warning, Error, Critical), PostgreSQL WARNING/PANIC/LOG and syslog notice/crit all became INFO. Matching is now case-insensitive and covers those names. Co-Authored-By: Claude Opus 5.5 --- output/otlp_grpc/metrics_traces.go | 80 ++++++---- output/otlp_grpc/metrics_traces_test.go | 8 +- output/otlp_grpc/otlp_grpc.go | 100 ++++++------ output/otlp_grpc/resource.go | 90 +++++++++++ output/otlp_grpc/resource_test.go | 198 ++++++++++++++++++++++++ output/otlp_grpc/send_test.go | 12 +- output/otlp_grpc/severity_test.go | 48 ++++++ 7 files changed, 446 insertions(+), 90 deletions(-) create mode 100644 output/otlp_grpc/resource.go create mode 100644 output/otlp_grpc/resource_test.go create mode 100644 output/otlp_grpc/severity_test.go diff --git a/output/otlp_grpc/metrics_traces.go b/output/otlp_grpc/metrics_traces.go index 61a506f3..04a41b0e 100644 --- a/output/otlp_grpc/metrics_traces.go +++ b/output/otlp_grpc/metrics_traces.go @@ -17,7 +17,7 @@ func (o *OTLPGrpc) WriteMetric(ctx context.Context, data output.MetricRecord) er pbMetric := convertMetricRecord(data) select { - case o.metricChan <- pbMetric: + case o.metricChan <- newEntry(pbMetric, data.Metadata.Resource): o.metrics.BlitzOutputEntriesReceivedCounter.Add(ctx, 1, outputType, "metrics") return nil case <-ctx.Done(): @@ -32,7 +32,7 @@ func (o *OTLPGrpc) WriteTrace(ctx context.Context, data output.TraceRecord) erro pbSpan := convertTraceRecord(data) select { - case o.traceChan <- pbSpan: + case o.traceChan <- newEntry(pbSpan, data.Metadata.Resource): o.metrics.BlitzOutputEntriesReceivedCounter.Add(ctx, 1, outputType, "traces") return nil case <-ctx.Done(): @@ -230,18 +230,24 @@ func hexByte(c byte) byte { } } -// buildMetricRequest builds an OTLP ExportMetricsServiceRequest from prepared metrics. -func buildMetricRequest(metrics []*metricspb.Metric, resource map[string]any) *metricspb.ResourceMetrics { - resourceAttrs := make([]*commonpb.KeyValue, 0, len(resource)+1) - resourceAttrs = append(resourceAttrs, &commonpb.KeyValue{ - Key: "service.name", - Value: &commonpb.AnyValue{Value: &commonpb.AnyValue_StringValue{StringValue: "blitz"}}, - }) - resourceAttrs = append(resourceAttrs, anyMapToKeyValues(resource)...) +// buildMetricRequests builds one ResourceMetrics per distinct resource in the +// batch, so each metric keeps the resource attributes (host.name, ...) its +// generator set in Metadata.Resource. +func buildMetricRequests(metrics []*entry[*metricspb.Metric]) []*metricspb.ResourceMetrics { + groups := groupByResource(metrics) + out := make([]*metricspb.ResourceMetrics, 0, len(groups)) + for _, g := range groups { + out = append(out, buildMetricRequest(g.items, g.resource)) + } + return out +} +// buildMetricRequest builds an OTLP ResourceMetrics from prepared metrics that +// share one resource. +func buildMetricRequest(metrics []*metricspb.Metric, resource map[string]any) *metricspb.ResourceMetrics { return &metricspb.ResourceMetrics{ Resource: &resourcepb.Resource{ - Attributes: resourceAttrs, + Attributes: resourceAttributes(resource), }, ScopeMetrics: []*metricspb.ScopeMetrics{ { @@ -251,16 +257,24 @@ func buildMetricRequest(metrics []*metricspb.Metric, resource map[string]any) *m } } -// buildTraceRequest builds an OTLP ResourceSpans from prepared spans. -func buildTraceRequest(spans []*tracepb.Span) *tracepb.ResourceSpans { +// buildTraceRequests builds one ResourceSpans per distinct resource in the +// batch, so each span keeps the resource attributes (host.name, service.name, +// ...) its generator set in Metadata.Resource. +func buildTraceRequests(spans []*entry[*tracepb.Span]) []*tracepb.ResourceSpans { + groups := groupByResource(spans) + out := make([]*tracepb.ResourceSpans, 0, len(groups)) + for _, g := range groups { + out = append(out, buildTraceRequest(g.items, g.resource)) + } + return out +} + +// buildTraceRequest builds an OTLP ResourceSpans from prepared spans that +// share one resource. +func buildTraceRequest(spans []*tracepb.Span, resource map[string]any) *tracepb.ResourceSpans { return &tracepb.ResourceSpans{ Resource: &resourcepb.Resource{ - Attributes: []*commonpb.KeyValue{ - { - Key: "service.name", - Value: &commonpb.AnyValue{Value: &commonpb.AnyValue_StringValue{StringValue: "blitz"}}, - }, - }, + Attributes: resourceAttributes(resource), }, ScopeSpans: []*tracepb.ScopeSpans{ { @@ -272,50 +286,50 @@ func buildTraceRequest(spans []*tracepb.Span) *tracepb.ResourceSpans { // metricBatch holds a batch of metrics to be sent type metricBatch struct { - metrics []*metricspb.Metric + metrics []*entry[*metricspb.Metric] maxSize int timer *time.Timer } func newMetricBatch(maxSize int, timeout time.Duration) *metricBatch { return &metricBatch{ - metrics: make([]*metricspb.Metric, 0, maxSize), + metrics: make([]*entry[*metricspb.Metric], 0, maxSize), maxSize: maxSize, timer: time.NewTimer(timeout), } } -func (b *metricBatch) add(m *metricspb.Metric) { b.metrics = append(b.metrics, m) } -func (b *metricBatch) isFull() bool { return len(b.metrics) >= b.maxSize } -func (b *metricBatch) isEmpty() bool { return len(b.metrics) == 0 } +func (b *metricBatch) add(m *entry[*metricspb.Metric]) { b.metrics = append(b.metrics, m) } +func (b *metricBatch) isFull() bool { return len(b.metrics) >= b.maxSize } +func (b *metricBatch) isEmpty() bool { return len(b.metrics) == 0 } -func (b *metricBatch) getAndClear() []*metricspb.Metric { +func (b *metricBatch) getAndClear() []*entry[*metricspb.Metric] { metrics := b.metrics - b.metrics = make([]*metricspb.Metric, 0, b.maxSize) + b.metrics = make([]*entry[*metricspb.Metric], 0, b.maxSize) return metrics } // traceBatch holds a batch of spans to be sent type traceBatch struct { - spans []*tracepb.Span + spans []*entry[*tracepb.Span] maxSize int timer *time.Timer } func newTraceBatch(maxSize int, timeout time.Duration) *traceBatch { return &traceBatch{ - spans: make([]*tracepb.Span, 0, maxSize), + spans: make([]*entry[*tracepb.Span], 0, maxSize), maxSize: maxSize, timer: time.NewTimer(timeout), } } -func (b *traceBatch) add(s *tracepb.Span) { b.spans = append(b.spans, s) } -func (b *traceBatch) isFull() bool { return len(b.spans) >= b.maxSize } -func (b *traceBatch) isEmpty() bool { return len(b.spans) == 0 } +func (b *traceBatch) add(s *entry[*tracepb.Span]) { b.spans = append(b.spans, s) } +func (b *traceBatch) isFull() bool { return len(b.spans) >= b.maxSize } +func (b *traceBatch) isEmpty() bool { return len(b.spans) == 0 } -func (b *traceBatch) getAndClear() []*tracepb.Span { +func (b *traceBatch) getAndClear() []*entry[*tracepb.Span] { spans := b.spans - b.spans = make([]*tracepb.Span, 0, b.maxSize) + b.spans = make([]*entry[*tracepb.Span], 0, b.maxSize) return spans } diff --git a/output/otlp_grpc/metrics_traces_test.go b/output/otlp_grpc/metrics_traces_test.go index 8f2b479c..a7bc32f5 100644 --- a/output/otlp_grpc/metrics_traces_test.go +++ b/output/otlp_grpc/metrics_traces_test.go @@ -167,11 +167,11 @@ func TestMetricBatch(t *testing.T) { assert.True(t, batch.isEmpty()) assert.False(t, batch.isFull()) - batch.add(&metricspb.Metric{Name: "test1"}) + batch.add(newEntry(&metricspb.Metric{Name: "test1"}, nil)) assert.False(t, batch.isEmpty()) assert.False(t, batch.isFull()) - batch.add(&metricspb.Metric{Name: "test2"}) + batch.add(newEntry(&metricspb.Metric{Name: "test2"}, nil)) assert.True(t, batch.isFull()) metrics := batch.getAndClear() @@ -184,8 +184,8 @@ func TestTraceBatch(t *testing.T) { batch := newTraceBatch(2, time.Second) assert.True(t, batch.isEmpty()) - batch.add(&tracepb.Span{Name: "span1"}) - batch.add(&tracepb.Span{Name: "span2"}) + batch.add(newEntry(&tracepb.Span{Name: "span1"}, nil)) + batch.add(newEntry(&tracepb.Span{Name: "span2"}, nil)) assert.True(t, batch.isFull()) spans := batch.getAndClear() diff --git a/output/otlp_grpc/otlp_grpc.go b/output/otlp_grpc/otlp_grpc.go index 4562def1..2f6a544d 100644 --- a/output/otlp_grpc/otlp_grpc.go +++ b/output/otlp_grpc/otlp_grpc.go @@ -5,6 +5,7 @@ import ( "crypto/tls" "fmt" "math" + "strings" "sync" "time" @@ -168,9 +169,9 @@ type OTLPGrpc struct { workers int insecure bool tlsConfig *tls.Config - dataChan chan *logspb.LogRecord - metricChan chan *metricspb.Metric - traceChan chan *tracepb.Span + dataChan chan *entry[*logspb.LogRecord] + metricChan chan *entry[*metricspb.Metric] + traceChan chan *entry[*tracepb.Span] ctx context.Context cancel context.CancelFunc workerManager *workermanager.WorkerManager @@ -241,9 +242,9 @@ func New(logger *zap.Logger, opts ...OTLPGrpcOption) (*OTLPGrpc, error) { workers: cfg.workers, insecure: cfg.insecure, tlsConfig: cfg.tlsConfig, - dataChan: make(chan *logspb.LogRecord, DefaultOTLPGrpcChannelSize), - metricChan: make(chan *metricspb.Metric, DefaultOTLPGrpcChannelSize), - traceChan: make(chan *tracepb.Span, DefaultOTLPGrpcChannelSize), + dataChan: make(chan *entry[*logspb.LogRecord], DefaultOTLPGrpcChannelSize), + metricChan: make(chan *entry[*metricspb.Metric], DefaultOTLPGrpcChannelSize), + traceChan: make(chan *entry[*tracepb.Span], DefaultOTLPGrpcChannelSize), ctx: ctx, cancel: cancel, batchTimeout: cfg.batchTimeout, @@ -349,9 +350,10 @@ func (o *OTLPGrpc) Write(ctx context.Context, data output.LogRecord) error { }, }, } + record.Attributes = append(record.Attributes, anyMapToKeyValues(data.Metadata.Attributes)...) select { - case o.dataChan <- record: + case o.dataChan <- newEntry(record, data.Metadata.Resource): o.metrics.BlitzOutputEntriesReceivedCounter.Add(ctx, 1, outputType, "logs") return nil case <-ctx.Done(): @@ -542,9 +544,8 @@ func (o *OTLPGrpc) sendMetricBatch(client collectormetrics.MetricsServiceClient, span.SetAttributes(attribute.Int("blitz.batch.size", len(metrics)), attribute.String("blitz.signal", "metrics")) defer span.End() - rm := buildMetricRequest(metrics, nil) request := &collectormetrics.ExportMetricsServiceRequest{ - ResourceMetrics: []*metricspb.ResourceMetrics{rm}, + ResourceMetrics: buildMetricRequests(metrics), } ctx, cancel := context.WithTimeout(o.ctx, o.batchTimeout) @@ -578,9 +579,8 @@ func (o *OTLPGrpc) sendTraceBatch(client collectortrace.TraceServiceClient, batc span.SetAttributes(attribute.Int("blitz.batch.size", len(spans)), attribute.String("blitz.signal", "traces")) defer span.End() - rs := buildTraceRequest(spans) request := &collectortrace.ExportTraceServiceRequest{ - ResourceSpans: []*tracepb.ResourceSpans{rs}, + ResourceSpans: buildTraceRequests(spans), } ctx, cancel := context.WithTimeout(o.ctx, o.batchTimeout) @@ -629,7 +629,7 @@ func (o *OTLPGrpc) connect() (*grpc.ClientConn, error) { // logBatch holds a batch of logs to be sent type logBatch struct { - logs []*logspb.LogRecord + logs []*entry[*logspb.LogRecord] maxSize int timer *time.Timer mu sync.Mutex @@ -638,14 +638,14 @@ type logBatch struct { // newLogBatch creates a new log batch func newLogBatch(maxSize int, timeout time.Duration) *logBatch { return &logBatch{ - logs: make([]*logspb.LogRecord, 0, maxSize), + logs: make([]*entry[*logspb.LogRecord], 0, maxSize), maxSize: maxSize, timer: time.NewTimer(timeout), } } // add adds a log to the batch -func (b *logBatch) add(data *logspb.LogRecord) { +func (b *logBatch) add(data *entry[*logspb.LogRecord]) { b.mu.Lock() defer b.mu.Unlock() @@ -667,11 +667,11 @@ func (b *logBatch) isEmpty() bool { } // getAndClear returns all logs and clears the batch -func (b *logBatch) getAndClear() []*logspb.LogRecord { +func (b *logBatch) getAndClear() []*entry[*logspb.LogRecord] { b.mu.Lock() defer b.mu.Unlock() logs := b.logs - b.logs = make([]*logspb.LogRecord, 0, b.maxSize) + b.logs = make([]*entry[*logspb.LogRecord], 0, b.maxSize) return logs } @@ -731,37 +731,31 @@ func (o *OTLPGrpc) flushBatch(client collectorlogs.LogsServiceClient, batch *log return o.sendBatch(client, batch) } -// buildOTLPRequest builds an OTLP ExportLogsServiceRequest from prepared LogRecord entries -func (o *OTLPGrpc) buildOTLPRequest(logs []*logspb.LogRecord) *collectorlogs.ExportLogsServiceRequest { - resourceLogs := &logspb.ResourceLogs{ - Resource: &resourcepb.Resource{ - Attributes: []*commonpb.KeyValue{ - { - Key: "service.name", - Value: &commonpb.AnyValue{ - Value: &commonpb.AnyValue_StringValue{ - StringValue: "blitz", - }, - }, - }, - }, - }, - ScopeLogs: []*logspb.ScopeLogs{ - { - LogRecords: make([]*logspb.LogRecord, 0, len(logs)), - }, - }, - } - - for _, logRecord := range logs { - if logRecord == nil { +// buildOTLPRequest builds an OTLP ExportLogsServiceRequest from prepared log +// entries, with one ResourceLogs per distinct resource so each record keeps +// the resource attributes (host.name, telemetry.source, ...) its generator +// set in Metadata.Resource. +func (o *OTLPGrpc) buildOTLPRequest(logs []*entry[*logspb.LogRecord]) *collectorlogs.ExportLogsServiceRequest { + groups := groupByResource(logs) + resourceLogs := make([]*logspb.ResourceLogs, 0, len(groups)) + for _, g := range groups { + records := make([]*logspb.LogRecord, 0, len(g.items)) + for _, r := range g.items { + if r != nil { + records = append(records, r) + } + } + if len(records) == 0 { continue } - resourceLogs.ScopeLogs[0].LogRecords = append(resourceLogs.ScopeLogs[0].LogRecords, logRecord) + resourceLogs = append(resourceLogs, &logspb.ResourceLogs{ + Resource: &resourcepb.Resource{Attributes: resourceAttributes(g.resource)}, + ScopeLogs: []*logspb.ScopeLogs{{LogRecords: records}}, + }) } return &collectorlogs.ExportLogsServiceRequest{ - ResourceLogs: []*logspb.ResourceLogs{resourceLogs}, + ResourceLogs: resourceLogs, } } @@ -849,19 +843,31 @@ func (o *OTLPGrpc) SupportedTelemetry() []telemetry.Type { return []telemetry.Type{telemetry.Logs, telemetry.Metrics, telemetry.Traces} } -// mapSeverityNumber maps string log levels to OTLP severity numbers +// mapSeverityNumber maps string log levels to OTLP severity numbers. +// Matching is case-insensitive and covers the level names generators emit: +// lowercase levels (kubernetes, apache error), Windows Event Log level names +// (Verbose, Information, Warning, Critical), syslog-style names (notice, crit) +// and PostgreSQL's LOG, WARNING and PANIC. Unknown levels map to INFO. func (o *OTLPGrpc) mapSeverityNumber(level string) logspb.SeverityNumber { - switch level { + switch strings.ToUpper(level) { + case "TRACE", "VERBOSE": + return logspb.SeverityNumber_SEVERITY_NUMBER_TRACE case "DEBUG": return logspb.SeverityNumber_SEVERITY_NUMBER_DEBUG - case "INFO": + case "INFO", "INFORMATION", "INFORMATIONAL", "LOG": return logspb.SeverityNumber_SEVERITY_NUMBER_INFO - case "WARN": + case "NOTICE": + return logspb.SeverityNumber_SEVERITY_NUMBER_INFO2 + case "WARN", "WARNING": return logspb.SeverityNumber_SEVERITY_NUMBER_WARN - case "ERROR": + case "ERROR", "ERR": return logspb.SeverityNumber_SEVERITY_NUMBER_ERROR + case "CRIT", "CRITICAL", "ALERT": + return logspb.SeverityNumber_SEVERITY_NUMBER_FATAL case "FATAL": return logspb.SeverityNumber_SEVERITY_NUMBER_FATAL2 + case "PANIC", "EMERG", "EMERGENCY": + return logspb.SeverityNumber_SEVERITY_NUMBER_FATAL4 default: return logspb.SeverityNumber_SEVERITY_NUMBER_INFO } diff --git a/output/otlp_grpc/resource.go b/output/otlp_grpc/resource.go new file mode 100644 index 00000000..5a0f5bcb --- /dev/null +++ b/output/otlp_grpc/resource.go @@ -0,0 +1,90 @@ +package otlpgrpc + +import ( + "fmt" + "sort" + "strings" + + commonpb "go.opentelemetry.io/proto/otlp/common/v1" +) + +// entry is a prepared OTLP item (log record, metric or span) plus the +// resource of the generator that emitted it. The resource travels with the +// item through the channel and batch so a send can split the batch into one +// Resource{Logs,Metrics,Spans} per distinct resource. +type entry[T any] struct { + item T + resource map[string]any + resourceKey string +} + +func newEntry[T any](item T, resource map[string]any) *entry[T] { + return &entry[T]{item: item, resource: resource, resourceKey: resourceKey(resource)} +} + +// resourceGroup is the items of a batch that share one resource. +type resourceGroup[T any] struct { + resource map[string]any + items []T +} + +// groupByResource splits entries into one group per distinct resource, in +// order of first appearance. nil entries are skipped. +func groupByResource[T any](entries []*entry[T]) []*resourceGroup[T] { + groups := make([]*resourceGroup[T], 0, 1) + byKey := make(map[string]*resourceGroup[T]) + for _, e := range entries { + if e == nil { + continue + } + g, ok := byKey[e.resourceKey] + if !ok { + g = &resourceGroup[T]{resource: e.resource} + byKey[e.resourceKey] = g + groups = append(groups, g) + } + g.items = append(g.items, e.item) + } + return groups +} + +// resourceAttributes returns the OTLP resource attributes for a generator +// resource: service.name=blitz unless the resource sets its own, followed by +// the resource's attributes in key order. +func resourceAttributes(resource map[string]any) []*commonpb.KeyValue { + attrs := make([]*commonpb.KeyValue, 0, len(resource)+1) + if _, ok := resource["service.name"]; !ok { + attrs = append(attrs, &commonpb.KeyValue{ + Key: "service.name", + Value: &commonpb.AnyValue{Value: &commonpb.AnyValue_StringValue{StringValue: "blitz"}}, + }) + } + for _, k := range sortedKeys(resource) { + if av := toAnyValueSimple(resource[k]); av != nil { + attrs = append(attrs, &commonpb.KeyValue{Key: k, Value: av}) + } + } + return attrs +} + +// resourceKey returns a stable identity for a resource map, used to group +// items that share a resource. +func resourceKey(resource map[string]any) string { + if len(resource) == 0 { + return "" + } + var b strings.Builder + for _, k := range sortedKeys(resource) { + fmt.Fprintf(&b, "%q=%#v;", k, resource[k]) + } + return b.String() +} + +func sortedKeys(m map[string]any) []string { + keys := make([]string, 0, len(m)) + for k := range m { + keys = append(keys, k) + } + sort.Strings(keys) + return keys +} diff --git a/output/otlp_grpc/resource_test.go b/output/otlp_grpc/resource_test.go new file mode 100644 index 00000000..2e58ce47 --- /dev/null +++ b/output/otlp_grpc/resource_test.go @@ -0,0 +1,198 @@ +package otlpgrpc + +import ( + "context" + "testing" + "time" + + "github.com/observiq/blitz/output" + commonpb "go.opentelemetry.io/proto/otlp/common/v1" + logspb "go.opentelemetry.io/proto/otlp/logs/v1" + metricspb "go.opentelemetry.io/proto/otlp/metrics/v1" + tracepb "go.opentelemetry.io/proto/otlp/trace/v1" + "go.uber.org/zap" +) + +func attrString(t *testing.T, attrs []*commonpb.KeyValue, key string) (string, bool) { + t.Helper() + for _, kv := range attrs { + if kv.Key == key { + sv, ok := kv.Value.Value.(*commonpb.AnyValue_StringValue) + if !ok { + t.Fatalf("attribute %q should be a string, got %T", key, kv.Value.Value) + } + return sv.StringValue, true + } + } + return "", false +} + +func assertResource(t *testing.T, attrs []*commonpb.KeyValue, want map[string]string) { + t.Helper() + for key, w := range want { + got, ok := attrString(t, attrs, key) + if !ok || got != w { + t.Errorf("resource %q = %q (present %v), want %q", key, got, ok, w) + } + } +} + +func newTestOutput(t *testing.T) *OTLPGrpc { + t.Helper() + o, err := New(zap.NewNop()) + if err != nil { + t.Fatalf("New: %v", err) + } + return o +} + +// A log record's Metadata.Resource must become the OTLP resource, so +// collectors can route on resource.attributes["telemetry.source"]. +func TestWrite_LogResourceReachesOTLPResource(t *testing.T) { + o := newTestOutput(t) + + err := o.Write(context.Background(), output.LogRecord{ + Message: "GET / 200", + Metadata: output.LogRecordMetadata{ + Timestamp: time.Now(), + Resource: map[string]any{"host.name": "web-01", "telemetry.source": "nginx"}, + Attributes: map[string]any{"http.method": "GET"}, + }, + }) + if err != nil { + t.Fatalf("Write: %v", err) + } + + req := o.buildOTLPRequest([]*entry[*logspb.LogRecord]{<-o.dataChan}) + if len(req.ResourceLogs) != 1 { + t.Fatalf("want 1 ResourceLogs, got %d", len(req.ResourceLogs)) + } + rl := req.ResourceLogs[0] + assertResource(t, rl.Resource.Attributes, map[string]string{ + "service.name": "blitz", + "host.name": "web-01", + "telemetry.source": "nginx", + }) + + if got, _ := attrString(t, rl.ScopeLogs[0].LogRecords[0].Attributes, "http.method"); got != "GET" { + t.Errorf("record attribute http.method = %q, want GET", got) + } +} + +// Records from different generators sharing one worker batch must not be +// merged under a single resource. +func TestBuildOTLPRequest_GroupsByResource(t *testing.T) { + o := &OTLPGrpc{} + entries := []*entry[*logspb.LogRecord]{ + newEntry(&logspb.LogRecord{}, map[string]any{"telemetry.source": "nginx"}), + newEntry(&logspb.LogRecord{}, map[string]any{"telemetry.source": "okta"}), + newEntry(&logspb.LogRecord{}, map[string]any{"telemetry.source": "nginx"}), + newEntry(&logspb.LogRecord{}, nil), + nil, + } + + req := o.buildOTLPRequest(entries) + if len(req.ResourceLogs) != 3 { + t.Fatalf("want 3 ResourceLogs (nginx, okta, none), got %d", len(req.ResourceLogs)) + } + + want := []struct { + source string + count int + }{{"nginx", 2}, {"okta", 1}, {"", 1}} + for i, w := range want { + rl := req.ResourceLogs[i] + if got, _ := attrString(t, rl.Resource.Attributes, "telemetry.source"); got != w.source { + t.Errorf("ResourceLogs[%d] telemetry.source = %q, want %q", i, got, w.source) + } + if n := len(rl.ScopeLogs[0].LogRecords); n != w.count { + t.Errorf("ResourceLogs[%d] has %d records, want %d", i, n, w.count) + } + } + + // A record with no resource still identifies itself as blitz. + if got, _ := attrString(t, req.ResourceLogs[2].Resource.Attributes, "service.name"); got != "blitz" { + t.Errorf("resource-less ResourceLogs service.name = %q, want blitz", got) + } +} + +// A metric's Metadata.Resource must become the OTLP resource, so host metrics +// arrive with host.name, and metrics from different resources stay apart. +func TestWriteMetric_ResourceReachesOTLPResource(t *testing.T) { + o := newTestOutput(t) + v := int64(1) + + for _, host := range []string{"host-a", "host-b", "host-a"} { + err := o.WriteMetric(context.Background(), output.MetricRecord{ + Name: "system.cpu.time", + Type: output.MetricTypeSum, + IntValue: &v, + Metadata: output.MetricPointMetadata{ + Timestamp: time.Now(), + Resource: map[string]any{"host.name": host, "telemetry.source": "hostmetrics"}, + }, + }) + if err != nil { + t.Fatalf("WriteMetric: %v", err) + } + } + + rms := buildMetricRequests([]*entry[*metricspb.Metric]{<-o.metricChan, <-o.metricChan, <-o.metricChan}) + if len(rms) != 2 { + t.Fatalf("want 2 ResourceMetrics (host-a, host-b), got %d", len(rms)) + } + assertResource(t, rms[0].Resource.Attributes, map[string]string{ + "service.name": "blitz", + "host.name": "host-a", + "telemetry.source": "hostmetrics", + }) + if n := len(rms[0].ScopeMetrics[0].Metrics); n != 2 { + t.Errorf("host-a ResourceMetrics has %d metrics, want 2", n) + } + assertResource(t, rms[1].Resource.Attributes, map[string]string{"host.name": "host-b"}) +} + +// A span's Metadata.Resource must become the OTLP resource, including a +// generator-provided service.name in place of the blitz default. +func TestWriteTrace_ResourceReachesOTLPResource(t *testing.T) { + o := newTestOutput(t) + + err := o.WriteTrace(context.Background(), output.TraceRecord{ + TraceID: "0102030405060708090a0b0c0d0e0f10", + SpanID: "0102030405060708", + Name: "GET /checkout", + StartTime: time.Now(), + EndTime: time.Now(), + Metadata: output.SpanMetadata{ + Resource: map[string]any{"service.name": "checkout", "host.name": "app-01"}, + }, + }) + if err != nil { + t.Fatalf("WriteTrace: %v", err) + } + + rss := buildTraceRequests([]*entry[*tracepb.Span]{<-o.traceChan}) + if len(rss) != 1 { + t.Fatalf("want 1 ResourceSpans, got %d", len(rss)) + } + assertResource(t, rss[0].Resource.Attributes, map[string]string{ + "service.name": "checkout", + "host.name": "app-01", + }) +} + +// A generator that sets its own service.name must not get a second one. +func TestResourceAttributes_KeepsGeneratorServiceName(t *testing.T) { + n := 0 + for _, kv := range resourceAttributes(map[string]any{"service.name": "checkout"}) { + if kv.Key == "service.name" { + n++ + if kv.Value.GetStringValue() != "checkout" { + t.Errorf("service.name = %q, want checkout", kv.Value.GetStringValue()) + } + } + } + if n != 1 { + t.Errorf("want exactly one service.name attribute, got %d", n) + } +} diff --git a/output/otlp_grpc/send_test.go b/output/otlp_grpc/send_test.go index 756f1c6d..1885ad69 100644 --- a/output/otlp_grpc/send_test.go +++ b/output/otlp_grpc/send_test.go @@ -58,26 +58,26 @@ func TestOTLPGrpc_sendBatchesEmitSpans(t *testing.T) { o := testOTLP(t) lb := newLogBatch(10, time.Second) - lb.add(&logspb.LogRecord{}) + lb.add(newEntry(&logspb.LogRecord{}, nil)) require.NoError(t, o.sendBatch(mockLogsClient{}, lb)) lbErr := newLogBatch(10, time.Second) - lbErr.add(&logspb.LogRecord{}) + lbErr.add(newEntry(&logspb.LogRecord{}, nil)) require.Error(t, o.sendBatch(mockLogsClient{err: errors.New("boom")}, lbErr)) mb := newMetricBatch(10, time.Second) - mb.add(&metricspb.Metric{}) + mb.add(newEntry(&metricspb.Metric{}, nil)) require.NoError(t, o.sendMetricBatch(mockMetricsClient{}, mb)) mbErr := newMetricBatch(10, time.Second) - mbErr.add(&metricspb.Metric{}) + mbErr.add(newEntry(&metricspb.Metric{}, nil)) require.Error(t, o.sendMetricBatch(mockMetricsClient{err: errors.New("boom")}, mbErr)) tb := newTraceBatch(10, time.Second) - tb.add(&tracepb.Span{}) + tb.add(newEntry(&tracepb.Span{}, nil)) require.NoError(t, o.sendTraceBatch(mockTraceClient{}, tb)) tbErr := newTraceBatch(10, time.Second) - tbErr.add(&tracepb.Span{}) + tbErr.add(newEntry(&tracepb.Span{}, nil)) require.Error(t, o.sendTraceBatch(mockTraceClient{err: errors.New("boom")}, tbErr)) } diff --git a/output/otlp_grpc/severity_test.go b/output/otlp_grpc/severity_test.go new file mode 100644 index 00000000..2c18b221 --- /dev/null +++ b/output/otlp_grpc/severity_test.go @@ -0,0 +1,48 @@ +package otlpgrpc + +import ( + "testing" + + logspb "go.opentelemetry.io/proto/otlp/logs/v1" +) + +// Generators emit levels in several spellings. Each must map to its real +// severity, not fall through to INFO. +func TestMapSeverityNumber(t *testing.T) { + o := &OTLPGrpc{} + cases := map[string]logspb.SeverityNumber{ + // uppercase (existing behavior) + "DEBUG": logspb.SeverityNumber_SEVERITY_NUMBER_DEBUG, + "INFO": logspb.SeverityNumber_SEVERITY_NUMBER_INFO, + "WARN": logspb.SeverityNumber_SEVERITY_NUMBER_WARN, + "ERROR": logspb.SeverityNumber_SEVERITY_NUMBER_ERROR, + "FATAL": logspb.SeverityNumber_SEVERITY_NUMBER_FATAL2, + // lowercase (kubernetes, apache error) + "debug": logspb.SeverityNumber_SEVERITY_NUMBER_DEBUG, + "info": logspb.SeverityNumber_SEVERITY_NUMBER_INFO, + "warn": logspb.SeverityNumber_SEVERITY_NUMBER_WARN, + "error": logspb.SeverityNumber_SEVERITY_NUMBER_ERROR, + "fatal": logspb.SeverityNumber_SEVERITY_NUMBER_FATAL2, + // Windows Event Log level names + "Verbose": logspb.SeverityNumber_SEVERITY_NUMBER_TRACE, + "Information": logspb.SeverityNumber_SEVERITY_NUMBER_INFO, + "Warning": logspb.SeverityNumber_SEVERITY_NUMBER_WARN, + "Error": logspb.SeverityNumber_SEVERITY_NUMBER_ERROR, + "Critical": logspb.SeverityNumber_SEVERITY_NUMBER_FATAL, + // syslog-style (apache error) + "notice": logspb.SeverityNumber_SEVERITY_NUMBER_INFO2, + "crit": logspb.SeverityNumber_SEVERITY_NUMBER_FATAL, + // PostgreSQL + "LOG": logspb.SeverityNumber_SEVERITY_NUMBER_INFO, + "WARNING": logspb.SeverityNumber_SEVERITY_NUMBER_WARN, + "PANIC": logspb.SeverityNumber_SEVERITY_NUMBER_FATAL4, + // unknown + "": logspb.SeverityNumber_SEVERITY_NUMBER_INFO, + "bogus": logspb.SeverityNumber_SEVERITY_NUMBER_INFO, + } + for level, want := range cases { + if got := o.mapSeverityNumber(level); got != want { + t.Errorf("mapSeverityNumber(%q) = %v, want %v", level, got, want) + } + } +} From e1358fcbdf51d3718e2bc1d1126d8bf49063d819 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adnan=20Rahi=C4=87?= Date: Wed, 7 Oct 2026 13:58:50 +0200 Subject: [PATCH 2/2] fix(hostmetrics): honor generator.hostmetrics.os when picking the fleet host Since generators draw their host identity from the simulated environment (PIPE-1036), hostmetrics took whatever system SystemForKey("hostmetrics") selected and used that system's OS, ignoring generator.hostmetrics.os. With the default randomized fleet that is often a Windows host, so BLITZ_GENERATOR_HOSTMETRICS_OS=linux (also the default) still produced Windows metrics. Add Environment.SystemForKeyWithOS, the same key-hashed selection restricted to systems running a given OS, and use it for hostmetrics with the configured OS (empty means linux). When no system runs that OS the generator falls back to its synthetic host built from OS and Hostname. Co-Authored-By: Claude Opus 5.5 --- internal/datagen/environment.go | 25 ++++++++++++-- internal/datagen/environment_selector_test.go | 33 +++++++++++++++++++ internal/dispatch/embed.go | 25 +++++++++++++- internal/dispatch/embed_test.go | 27 +++++++++++++++ 4 files changed, 106 insertions(+), 4 deletions(-) diff --git a/internal/datagen/environment.go b/internal/datagen/environment.go index 67414200..50347917 100644 --- a/internal/datagen/environment.go +++ b/internal/datagen/environment.go @@ -43,7 +43,26 @@ func (e *Environment) AllNetworkSystems() []*NetworkSystemIdentity { return e.Ne // identity once and attributes every record it emits consistently. Returns nil // when the environment has no systems. func (e *Environment) SystemForKey(key string) *SystemIdentity { - if len(e.Systems) == 0 { + return systemForKey(e.Systems, key) +} + +// SystemForKeyWithOS is SystemForKey restricted to the environment's systems +// running os: the same key always maps to the same system of that OS. It is +// for generators whose output depends on the host OS (hostmetrics), so a +// configured OS is honored instead of inheriting whatever OS the key-selected +// system happens to run. Returns nil when no system runs os. +func (e *Environment) SystemForKeyWithOS(key string, os OSType) *SystemIdentity { + matching := make([]*SystemIdentity, 0, len(e.Systems)) + for _, s := range e.Systems { + if s != nil && s.OSInfo.Type == os { + matching = append(matching, s) + } + } + return systemForKey(matching, key) +} + +func systemForKey(systems []*SystemIdentity, key string) *SystemIdentity { + if len(systems) == 0 { return nil } h := fnv.New32a() @@ -51,8 +70,8 @@ func (e *Environment) SystemForKey(key string) *SystemIdentity { // int64 throughout: uint32->int64 and int->int64 are widening (never // negative, never truncating), so this is correct on 32-bit targets and // avoids an int->uint32 narrowing conversion. - idx := int64(h.Sum32()) % int64(len(e.Systems)) - return e.Systems[idx] + idx := int64(h.Sum32()) % int64(len(systems)) + return systems[idx] } // EnvironmentOpts controls the size and shape of the generated environment. diff --git a/internal/datagen/environment_selector_test.go b/internal/datagen/environment_selector_test.go index 1af5d8d7..39b49305 100644 --- a/internal/datagen/environment_selector_test.go +++ b/internal/datagen/environment_selector_test.go @@ -41,3 +41,36 @@ func TestEnvironmentSystemForKey(t *testing.T) { t.Error("SystemForKey on an empty environment should return nil") } } + +func TestEnvironmentSystemForKeyWithOS(t *testing.T) { + env := &Environment{Systems: []*SystemIdentity{ + {Hostname: "win-1", OSInfo: OSInfo{Type: OSWindows}}, + {Hostname: "lin-1", OSInfo: OSInfo{Type: OSLinux}}, + {Hostname: "win-2", OSInfo: OSInfo{Type: OSWindows}}, + {Hostname: "lin-2", OSInfo: OSInfo{Type: OSLinux}}, + nil, + }} + + // Only systems running the requested OS are eligible, for every key. + for _, k := range []string{"hostmetrics", "apache", "nginx", "postgres", "wel", "traces"} { + if s := env.SystemForKeyWithOS(k, OSLinux); s == nil || s.OSInfo.Type != OSLinux { + t.Errorf("SystemForKeyWithOS(%q, linux) = %+v, want a linux system", k, s) + } + if s := env.SystemForKeyWithOS(k, OSWindows); s == nil || s.OSInfo.Type != OSWindows { + t.Errorf("SystemForKeyWithOS(%q, windows) = %+v, want a windows system", k, s) + } + } + + // Deterministic for a repeated key. + first := env.SystemForKeyWithOS("hostmetrics", OSLinux) + for i := 0; i < 5; i++ { + if env.SystemForKeyWithOS("hostmetrics", OSLinux) != first { + t.Fatal("SystemForKeyWithOS is not deterministic for a repeated key") + } + } + + // No system runs the OS: nil, not a system of another OS. + if s := env.SystemForKeyWithOS("hostmetrics", OSMacOS); s != nil { + t.Errorf("SystemForKeyWithOS(macos) = %+v, want nil", s) + } +} diff --git a/internal/dispatch/embed.go b/internal/dispatch/embed.go index d6844f5b..ee1b4d67 100644 --- a/internal/dispatch/embed.go +++ b/internal/dispatch/embed.go @@ -185,7 +185,7 @@ func ForEmbed(logger *zap.Logger, genCfg config.Generator, consumers EmbedConsum ScraperNames: genCfg.HostMetrics.Scrapers, Consumer: consumers.MetricConsumer, Seed: yamlSeedDefault(genCfg.HostMetrics.Seed), - Identity: hostIdentity(env, genCfg.Type), + Identity: hostMetricsIdentity(env, genCfg.HostMetrics.OS), Telemetry: tel, }) case config.GeneratorTypeTraces: @@ -249,6 +249,29 @@ func hostIdentity(env *datagen.Environment, component config.GeneratorType) *dat return env.SystemForKey(string(component)) } +// hostMetricsIdentity resolves the simulated host for the hostmetrics +// generator. Its metrics are OS-specific and its OS is configurable +// (generator.hostmetrics.os), so the host is chosen only among the +// environment's systems running that OS; otherwise the key-selected system's +// OS silently overrode the setting. An empty OS means linux, the config +// default. Returns nil (the generator then synthesizes a host from OS and +// Hostname) when no environment is configured, the OS is not recognized, or +// no system in the environment runs it. +func hostMetricsIdentity(env *datagen.Environment, osName string) *datagen.SystemIdentity { + if env == nil { + return nil + } + osType := datagen.OSLinux + if osName != "" { + parsed, err := datagen.ParseOSType(osName) + if err != nil { + return nil + } + osType = parsed + } + return env.SystemForKeyWithOS(string(config.GeneratorTypeHostMetrics), osType) +} + // yamlSeedDefault translates a YAML-loaded Seed value into the // generator-Config Seed value, applying the "stochastic by default" // architectural intent for YAML users. YAML zero-value (omitted `seed:` diff --git a/internal/dispatch/embed_test.go b/internal/dispatch/embed_test.go index ba0bae62..3ac89e78 100644 --- a/internal/dispatch/embed_test.go +++ b/internal/dispatch/embed_test.go @@ -133,6 +133,33 @@ func TestHostIdentityResolvesFromEnvironment(t *testing.T) { assert.Equal(t, "PANTHEON-01", got.Hostname) } +// TestHostMetricsIdentityHonorsConfiguredOS covers the hostmetrics-specific +// resolution: the selected host always runs the configured OS (empty means +// linux), and an OS no system runs falls back to nil (synthetic host). +func TestHostMetricsIdentityHonorsConfiguredOS(t *testing.T) { + assert.Nil(t, hostMetricsIdentity(nil, "linux")) + + env := &datagen.Environment{ + Systems: []*datagen.SystemIdentity{ + {Hostname: "FLORA-WORKER09", OSInfo: datagen.OSInfo{Type: datagen.OSWindows}}, + {Hostname: "mimir-db-14", OSInfo: datagen.OSInfo{Type: datagen.OSLinux}}, + }, + } + + for _, osName := range []string{"linux", "LINUX", ""} { + got := hostMetricsIdentity(env, osName) + require.NotNil(t, got, "os %q", osName) + assert.Equal(t, "mimir-db-14", got.Hostname, "os %q", osName) + } + + got := hostMetricsIdentity(env, "windows") + require.NotNil(t, got) + assert.Equal(t, "FLORA-WORKER09", got.Hostname) + + assert.Nil(t, hostMetricsIdentity(env, "macos"), "no macos system: fall back to a synthetic host") + assert.Nil(t, hostMetricsIdentity(env, "solaris"), "unknown OS: fall back to a synthetic host") +} + // TestForEmbedHostMetricsWiresEnvironmentIdentity proves the full wiring: a // hostmetrics module built through ForEmbed with an environment emits points // carrying the resolved simulated host's identity attributes.