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
25 changes: 22 additions & 3 deletions internal/datagen/environment.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,16 +43,35 @@ 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()
_, _ = h.Write([]byte(key))
// 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.
Expand Down
33 changes: 33 additions & 0 deletions internal/datagen/environment_selector_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
}
25 changes: 24 additions & 1 deletion internal/dispatch/embed.go
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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:`
Expand Down
27 changes: 27 additions & 0 deletions internal/dispatch/embed_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
80 changes: 47 additions & 33 deletions output/otlp_grpc/metrics_traces.go
Original file line number Diff line number Diff line change
Expand Up @@ -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():
Expand All @@ -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():
Expand Down Expand Up @@ -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{
{
Expand All @@ -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{
{
Expand All @@ -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
}
8 changes: 4 additions & 4 deletions output/otlp_grpc/metrics_traces_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand All @@ -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()
Expand Down
Loading
Loading