diff --git a/cmd/blitz/main.go b/cmd/blitz/main.go index edf95ecc..847c2211 100644 --- a/cmd/blitz/main.go +++ b/cmd/blitz/main.go @@ -167,6 +167,7 @@ func run(cmd *cobra.Command, args []string) error { // Emit Warn-level banners for any deprecated generator types // configured by the user. Fires once per startup, not per record. config.LogGeneratorDeprecations(logger, cfg) + config.LogRemovedSettings(logger, cfg) if err := setupMetrics(ctx, cfg, logger); err != nil { logger.Error("Failed to setup metrics", zap.Error(err)) diff --git a/config/loader_test.go b/config/loader_test.go index b17e72d0..10fae044 100644 --- a/config/loader_test.go +++ b/config/loader_test.go @@ -168,7 +168,6 @@ func TestLoadModules_HostMetricsRequiresMetricConsumer(t *testing.T) { generator: type: hostmetrics hostmetrics: - workers: 1 rate: 1s os: linux output: diff --git a/docker/docker-compose.telemetry-generator.yml b/docker/docker-compose.telemetry-generator.yml index 7780a7e2..ba3d03b3 100644 --- a/docker/docker-compose.telemetry-generator.yml +++ b/docker/docker-compose.telemetry-generator.yml @@ -188,7 +188,6 @@ services: <<: *blitz-common environment: BLITZ_GENERATOR_TYPE: hostmetrics - BLITZ_GENERATOR_HOSTMETRICS_WORKERS: ${BLITZ_WORKERS:-1} BLITZ_GENERATOR_HOSTMETRICS_RATE: ${BLITZ_RATE:-1s} BLITZ_GENERATOR_HOSTMETRICS_OS: linux BLITZ_OUTPUT_TYPE: otlp-grpc @@ -200,7 +199,6 @@ services: <<: *blitz-common environment: BLITZ_GENERATOR_TYPE: hostmetrics - BLITZ_GENERATOR_HOSTMETRICS_WORKERS: ${BLITZ_WORKERS:-1} BLITZ_GENERATOR_HOSTMETRICS_RATE: ${BLITZ_RATE:-1s} BLITZ_GENERATOR_HOSTMETRICS_OS: windows BLITZ_OUTPUT_TYPE: otlp-grpc diff --git a/docs/generator/hostmetrics.md b/docs/generator/hostmetrics.md index 552b3436..6c2ebef7 100644 --- a/docs/generator/hostmetrics.md +++ b/docs/generator/hostmetrics.md @@ -28,19 +28,21 @@ The generator includes 8 scrapers, each producing metrics for a specific subsyst | YAML Path | Flag Name | Environment Variable | Default | Description | |------------------------------------|------------------------------------|-----------------------------------------|-----------|----------------------------------------------------------------------| | `generator.type` | `--generator-type` | `BLITZ_GENERATOR_TYPE` | `nop` | Generator type. Set to `hostmetrics` to use this generator. | -| `generator.hostmetrics.workers` | `--generator-hostmetrics-workers` | `BLITZ_GENERATOR_HOSTMETRICS_WORKERS` | `1` | Number of worker goroutines. | | `generator.hostmetrics.rate` | `--generator-hostmetrics-rate` | `BLITZ_GENERATOR_HOSTMETRICS_RATE` | `1s` | Scrape interval for host metrics. | | `generator.hostmetrics.os` | `--generator-hostmetrics-os` | `BLITZ_GENERATOR_HOSTMETRICS_OS` | `linux` | Simulated operating system. One of: `linux`, `windows`. | | `generator.hostmetrics.hostname` | `--generator-hostmetrics-hostname` | `BLITZ_GENERATOR_HOSTMETRICS_HOSTNAME` | (random) | Simulated hostname. If empty, a random hostname is generated. | | `generator.hostmetrics.scrapers` | `--generator-hostmetrics-scrapers` | `BLITZ_GENERATOR_HOSTMETRICS_SCRAPERS` | (all) | Scrapers to enable. If empty, all scrapers are enabled. | +`generator.hostmetrics.workers` has been removed. One simulated host runs one worker, so use `rate` for more +frequent writes and add `generators:` entries for more hosts. A configured value is ignored with a warning, and is +expected to fail config validation as of v0.25.0. + ## Example Configuration ```yaml generator: type: hostmetrics hostmetrics: - workers: 1 rate: 1s os: linux scrapers: diff --git a/embed/record.go b/embed/record.go index ee86e5e0..f34f4253 100644 --- a/embed/record.go +++ b/embed/record.go @@ -45,13 +45,13 @@ type LogRecordMetadata struct { type MetricType string const ( - // MetricTypeGauge represents a gauge metric. + // MetricTypeGauge is a point-in-time value that can rise or fall. MetricTypeGauge MetricType = "gauge" - // MetricTypeSum represents a sum metric. + // MetricTypeSum is a non-monotonic sum (OTel Sum monotonic=false; Prometheus gauge). MetricTypeSum MetricType = "sum" - // MetricTypeCounter represents a counter metric. + // MetricTypeCounter is a monotonic cumulative count (OTel Sum monotonic=true; Prometheus counter, _total). MetricTypeCounter MetricType = "counter" - // MetricTypeHistogram represents a histogram metric. + // MetricTypeHistogram is a bucketed distribution (Prometheus _bucket/_sum/_count). MetricTypeHistogram MetricType = "histogram" ) diff --git a/generator/hostmetrics/hostmetrics.go b/generator/hostmetrics/hostmetrics.go index 4ab7b974..54764c35 100644 --- a/generator/hostmetrics/hostmetrics.go +++ b/generator/hostmetrics/hostmetrics.go @@ -24,9 +24,13 @@ const generatorType = "hostmetrics" type Config struct { // Logger is the zap logger used for diagnostic output. Required. Logger *zap.Logger - // Workers is the number of worker goroutines. Required, >= 1. + // Workers is ignored: one simulated host runs one worker, and Rate is the + // load knob. Parallel workers only emitted duplicate series for the same + // host. + // + // Deprecated: ignored; kept so existing embed callers still compile. Workers int - // Rate is the scrape interval per worker. Required, > 0. + // Rate is the scrape interval. Required, > 0. Rate time.Duration // OS is the simulated operating system ("linux" or "windows"). Ignored // when Identity is set (the identity's own OS is used instead). @@ -65,7 +69,6 @@ type Generator struct { embed.ProducerMarker logger *zap.Logger - workers int rate time.Duration osType string hostname string @@ -93,9 +96,6 @@ func New(cfg Config) (*Generator, error) { if cfg.Consumer == nil { return nil, fmt.Errorf("MetricConsumer cannot be nil") } - if cfg.Workers < 1 { - return nil, fmt.Errorf("workers must be 1 or greater, got %d", cfg.Workers) - } if cfg.Rate <= 0 { return nil, fmt.Errorf("rate must be greater than 0, got %s", cfg.Rate) } @@ -116,7 +116,6 @@ func New(cfg Config) (*Generator, error) { return &Generator{ logger: cfg.Logger.Named("generator-hostmetrics"), - workers: cfg.Workers, rate: cfg.Rate, osType: sys.OSInfo.Type.SemconvOSType(), hostname: sys.Hostname, @@ -169,27 +168,25 @@ func (g *Generator) SetCountTracker(tracker *count.Tracker) { g.tracker = tracker } -// Start launches the worker goroutines. +// Start launches the single worker: one simulated host emits one sample per +// series per rate interval. func (g *Generator) Start(_ context.Context) error { g.logger.Info("Starting host metrics generator", - zap.Int("workers", g.workers), zap.Duration("rate", g.rate), zap.String("os.type", g.osType), zap.String("hostname", g.hostname), zap.Int("scrapers", len(g.scrapers)), ) - g.metrics.BlitzGeneratorActiveWorkersGauge.Record(context.Background(), int64(g.workers), generatorType) + g.metrics.BlitzGeneratorActiveWorkersGauge.Record(context.Background(), 1, generatorType) - for i := range g.workers { - g.wg.Add(1) - go g.worker(i) - } + g.wg.Add(1) + go g.worker(0) return nil } -// Stop signals workers to drain and waits for them to exit. +// Stop signals the worker to drain and waits for it to exit. func (g *Generator) Stop(ctx context.Context) error { g.logger.Info("Stopping host metrics generator") @@ -253,9 +250,8 @@ func (g *Generator) scrape(r *rand.Rand) { // The host-identity resource is fixed for this generator's lifetime, so it // is built once (StaticResources in New) and shared read-only across every - // scrape and worker. Scrapers only attach it to the MetricRecords they - // return — they never mutate it — so handing out the zero-allocation shared - // map is safe under concurrent workers. + // scrape. Scrapers only attach it to the MetricRecords they return — they + // never mutate it — so handing out the zero-allocation shared map is safe. res := g.static.Record() for _, scraper := range g.scrapers { diff --git a/generator/hostmetrics/hostmetrics_test.go b/generator/hostmetrics/hostmetrics_test.go index 8e170340..3e19c0ae 100644 --- a/generator/hostmetrics/hostmetrics_test.go +++ b/generator/hostmetrics/hostmetrics_test.go @@ -13,6 +13,8 @@ import ( "github.com/observiq/blitz/telemetry" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "go.opentelemetry.io/otel/sdk/metric/metricdata" "go.uber.org/zap/zaptest" ) @@ -101,14 +103,45 @@ func TestNew(t *testing.T) { require.Error(t, err) }) - t.Run("invalid workers", func(t *testing.T) { + t.Run("workers unset is valid", func(t *testing.T) { cfg := baseCfg(t, &mockMetricConsumer{}) cfg.Workers = 0 _, err := New(cfg) - require.Error(t, err) + require.NoError(t, err) }) } +// TestWorkersIgnored asserts one simulated host runs exactly one worker no +// matter what Workers says: parallel workers would emit duplicate series for +// the same host, and rate is the load knob. +func TestWorkersIgnored(t *testing.T) { + reader := sdkmetric.NewManualReader() + mp := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader)) + + cfg := baseCfg(t, &mockMetricConsumer{}) + cfg.Workers = 4 + cfg.Telemetry = embed.TelemetrySettings{MeterProvider: mp} + g, err := New(cfg) + require.NoError(t, err) + require.NoError(t, g.Start(context.Background())) + t.Cleanup(func() { _ = g.Stop(context.Background()) }) + + var rm metricdata.ResourceMetrics + require.NoError(t, reader.Collect(context.Background(), &rm)) + var active int64 = -1 + for _, sm := range rm.ScopeMetrics { + for _, md := range sm.Metrics { + if md.Name != "blitz.generator.active_workers" { + continue + } + if gauge, ok := md.Data.(metricdata.Gauge[int64]); ok && len(gauge.DataPoints) > 0 { + active = gauge.DataPoints[0].Value + } + } + } + require.Equal(t, int64(1), active) +} + // TestNewProjectsIdentityResource confirms that when a resolved datagen // identity is supplied, the generator's static resource carries the full // host.* / os.* / deployment.* projection (os.type as the semconv value, so diff --git a/internal/config/generator_hostmetrics.go b/internal/config/generator_hostmetrics.go index ee2c0d4d..4c4510b6 100644 --- a/internal/config/generator_hostmetrics.go +++ b/internal/config/generator_hostmetrics.go @@ -9,7 +9,11 @@ import ( // HostMetricsGeneratorConfig contains configuration for host metrics generator type HostMetricsGeneratorConfig struct { - // Workers is the number of worker goroutines for host metrics generation + // Workers is removed and ignored: one simulated host runs one worker, and + // Rate is the load knob. Kept only so a still-configured value can be + // detected and warned about (see LogRemovedSettings). + // + // Deprecated: ignored; expected to fail validation as of v0.25.0. Workers int `yaml:"workers,omitempty" mapstructure:"workers,omitempty"` // Rate is the scrape interval for host metrics Rate time.Duration `yaml:"rate,omitempty" mapstructure:"rate,omitempty"` @@ -45,10 +49,6 @@ var ValidScrapers = []string{ // Validate validates the host metrics generator configuration func (c *HostMetricsGeneratorConfig) Validate() error { - if c.Workers < 1 { - return fmt.Errorf("hostmetrics generator workers must be 1 or greater, got %d", c.Workers) - } - if c.Rate <= 0 { return fmt.Errorf("hostmetrics generator rate must be positive, got %v", c.Rate) } diff --git a/internal/config/generator_hostmetrics_test.go b/internal/config/generator_hostmetrics_test.go index 7106d583..f5b3047c 100644 --- a/internal/config/generator_hostmetrics_test.go +++ b/internal/config/generator_hostmetrics_test.go @@ -39,13 +39,11 @@ func TestHostMetricsGeneratorConfig_Validate(t *testing.T) { }, }, { - name: "invalid workers", + // workers was removed (rate is the load knob); unset is valid. + name: "workers unset", config: HostMetricsGeneratorConfig{ - Workers: 0, - Rate: time.Second, + Rate: time.Second, }, - wantErr: true, - errMsg: "workers must be 1 or greater", }, { name: "invalid rate", diff --git a/internal/config/migrate.go b/internal/config/migrate.go index efae2894..e752c9d5 100644 --- a/internal/config/migrate.go +++ b/internal/config/migrate.go @@ -78,3 +78,32 @@ func LogGeneratorDeprecations(logger *zap.Logger, cfg *Config) { } } } + +// LogRemovedSettings emits a Warn once per startup for every configured +// setting that has been removed. A removed setting is ignored rather than +// rejected during a deprecation window, so existing configs keep loading. +// +// Currently emits warnings for: +// - generator.hostmetrics.workers: one simulated host runs one worker, and +// rate is the load knob. Parallel workers only duplicated the same host's +// series. +// +// TODO: hostmetrics `workers` was removed (rate is the load knob). Expected in +// v0.25.0: turn this warning into a config validation error and drop the +// deprecated --generator-hostmetrics-workers flag. +func LogRemovedSettings(logger *zap.Logger, cfg *Config) { + if logger == nil || cfg == nil { + return + } + for _, g := range cfg.EffectiveGenerators() { + if g.Type == GeneratorTypeHostMetrics && g.HostMetrics.Workers != 0 { + logger.Warn(HostMetricsWorkersRemoved) + } + } +} + +// HostMetricsWorkersRemoved is the warning for the removed +// generator.hostmetrics.workers setting, shared by the startup log and the +// deprecated CLI flag. +const HostMetricsWorkersRemoved = "`generator.hostmetrics.workers` is no longer supported and is ignored; " + + "use `rate` for more frequent writes. Setting it is expected to fail config validation as of v0.25.0." diff --git a/internal/config/migrate_test.go b/internal/config/migrate_test.go index b967bb0b..7b565bce 100644 --- a/internal/config/migrate_test.go +++ b/internal/config/migrate_test.go @@ -221,3 +221,62 @@ func TestLogGeneratorDeprecations_NilSafe(t *testing.T) { config.LogGeneratorDeprecations(nil, &config.Config{}) }) } + +// TestLogRemovedSettings_HostMetricsWorkers asserts a Warn fires when the +// removed generator.hostmetrics.workers setting is still configured, and that +// it points to rate and the expected v0.25.0 validation failure. +func TestLogRemovedSettings_HostMetricsWorkers(t *testing.T) { + core, recorded := observer.New(zap.WarnLevel) + logger := zap.New(core) + + cfg := &config.Config{ + Generator: config.Generator{ + Type: config.GeneratorTypeHostMetrics, + HostMetrics: config.HostMetricsGeneratorConfig{Workers: 4, Rate: time.Second}, + }, + } + + config.LogRemovedSettings(logger, cfg) + + entries := recorded.FilterMessageSnippet("generator.hostmetrics.workers").All() + require.Len(t, entries, 1) + assert.Contains(t, entries[0].Message, "ignored") + assert.Contains(t, entries[0].Message, "rate") + assert.Contains(t, entries[0].Message, "v0.25.0") +} + +// TestLogRemovedSettings_GeneratorsList covers hostmetrics entries in the +// generators: list, one Warn per offending entry. +func TestLogRemovedSettings_GeneratorsList(t *testing.T) { + core, recorded := observer.New(zap.WarnLevel) + logger := zap.New(core) + + cfg := &config.Config{ + Generators: []config.Generator{ + {Type: config.GeneratorTypeHostMetrics, HostMetrics: config.HostMetricsGeneratorConfig{Workers: 2, Rate: time.Second}}, + {Type: config.GeneratorTypeHostMetrics, HostMetrics: config.HostMetricsGeneratorConfig{Rate: time.Second}}, + }, + } + + config.LogRemovedSettings(logger, cfg) + + assert.Len(t, recorded.FilterMessageSnippet("generator.hostmetrics.workers").All(), 1) +} + +// TestLogRemovedSettings_NotSetNoWarn confirms no warning when workers is unset. +func TestLogRemovedSettings_NotSetNoWarn(t *testing.T) { + core, recorded := observer.New(zap.WarnLevel) + logger := zap.New(core) + + cfg := &config.Config{ + Generator: config.Generator{ + Type: config.GeneratorTypeHostMetrics, + HostMetrics: config.HostMetricsGeneratorConfig{Rate: time.Second}, + }, + } + + config.LogRemovedSettings(logger, cfg) + + assert.Empty(t, recorded.All()) + assert.NotPanics(t, func() { config.LogRemovedSettings(nil, nil) }) +} diff --git a/internal/config/override.go b/internal/config/override.go index e0594187..38f804fd 100644 --- a/internal/config/override.go +++ b/internal/config/override.go @@ -21,6 +21,9 @@ type Override struct { Usage string // Default is the default value for the override Default any + // Deprecated, when set, hides the flag and prints this message when it + // is used, while keeping it bound for a deprecation window. + Deprecated string } // NewOverride creates a new override @@ -37,6 +40,11 @@ func NewOverride(field, usage string, def any) *Override { // Bind binds the override to the viper instance func (o *Override) Bind(flags *pflag.FlagSet) error { flag := o.createFlag(flags) + if o.Deprecated != "" { + if err := flags.MarkDeprecated(o.Flag, o.Deprecated); err != nil { + return err + } + } if err := viper.BindPFlag(o.Field, flag); err != nil { return err } @@ -277,7 +285,15 @@ func DefaultOverrides() []*Override { NewOverride("generator.filegen.cache-ttl", "file cache time-to-live (0 = never expire)", time.Duration(0)), NewOverride("generator.okta.workers", "number of Okta generator workers", 1), NewOverride("generator.okta.rate", "rate at which Okta logs are generated per worker", 1*time.Second), - NewOverride("generator.hostmetrics.workers", "number of host metrics generator workers", 1), + // Removed setting kept bound for the deprecation window; see LogRemovedSettings. + &Override{ + Field: "generator.hostmetrics.workers", + Flag: createFlagName("generator.hostmetrics.workers"), + Env: createEnvName("generator.hostmetrics.workers"), + Usage: "removed: ignored, use --generator-hostmetrics-rate", + Default: 0, + Deprecated: HostMetricsWorkersRemoved, + }, NewOverride("generator.hostmetrics.rate", "scrape interval for host metrics generation", 1*time.Second), NewOverride("generator.hostmetrics.os", "simulated operating system. One of: linux|windows", "linux"), NewOverride("generator.hostmetrics.hostname", "simulated hostname (empty = random)", ""), diff --git a/internal/config/override_test.go b/internal/config/override_test.go index a81ae05c..5c824ec9 100644 --- a/internal/config/override_test.go +++ b/internal/config/override_test.go @@ -369,7 +369,6 @@ func TestOverrideDefaults(t *testing.T) { Rate: 1 * time.Second, }, HostMetrics: HostMetricsGeneratorConfig{ - Workers: 1, Rate: 1 * time.Second, OS: "linux", Scrapers: []string{}, @@ -997,3 +996,17 @@ func TestOverrideCoverage(t *testing.T) { t.Errorf("%s", report.String()) } } + +// TestHostMetricsWorkersFlagDeprecated asserts the removed +// --generator-hostmetrics-workers flag stays registered but is marked +// deprecated for the warning window, so CLI users see a warning instead of an +// unknown-flag failure. +func TestHostMetricsWorkersFlagDeprecated(t *testing.T) { + flagSet := pflag.NewFlagSet("test", pflag.ContinueOnError) + for _, o := range DefaultOverrides() { + require.NoError(t, o.Bind(flagSet)) + } + f := flagSet.Lookup("generator-hostmetrics-workers") + require.NotNil(t, f) + require.Contains(t, f.Deprecated, "rate") +} diff --git a/internal/dispatch/embed.go b/internal/dispatch/embed.go index d6844f5b..8e82686d 100644 --- a/internal/dispatch/embed.go +++ b/internal/dispatch/embed.go @@ -178,7 +178,6 @@ func ForEmbed(logger *zap.Logger, genCfg config.Generator, consumers EmbedConsum } return hostmetrics.New(hostmetrics.Config{ Logger: logger, - Workers: genCfg.HostMetrics.Workers, Rate: genCfg.HostMetrics.Rate, OS: genCfg.HostMetrics.OS, Hostname: genCfg.HostMetrics.Hostname, diff --git a/internal/prommap/prommap.go b/internal/prommap/prommap.go new file mode 100644 index 00000000..a7783959 --- /dev/null +++ b/internal/prommap/prommap.go @@ -0,0 +1,304 @@ +// Package prommap maps blitz embed MetricPoints to a neutral in-memory +// Prometheus model that both the scrape (text) and remote-write (protobuf) +// output encoders serialize from. It owns type mapping, name/label +// sanitization, suffixing, cumulative histogram expansion, and HELP/TYPE/unit. +package prommap + +import ( + "encoding/json" + "fmt" + "sort" + "strconv" + + "github.com/observiq/blitz/embed" +) + +// Type is a Prometheus metric family type. Counter and Histogram map directly; +// Sum maps to gauge (no monotonicity flag, so treated as non-monotonic). +type Type string + +const ( + // TypeGauge is a Prometheus gauge. + TypeGauge Type = "gauge" + // TypeCounter is a Prometheus counter. + TypeCounter Type = "counter" + // TypeHistogram is a Prometheus histogram. + TypeHistogram Type = "histogram" +) + +// Label is a single Prometheus label name/value pair. +type Label struct { + Name string + Value string +} + +// Sample is one time-series point. Name already carries any +// _total/_bucket/_sum/_count suffix. +type Sample struct { + Name string + Labels []Label + Value float64 + TimestampMS int64 +} + +// MetricFamily groups samples sharing a base name, type, and metadata. Name is +// the HELP/TYPE-line name (for a counter it already includes _total). +type MetricFamily struct { + Name string + Type Type + Help string + Unit string + Samples []Sample +} + +// Map projects a single embed MetricPoint into a MetricFamily. A gauge, sum, or +// counter yields one sample; a histogram yields cumulative _bucket series (with +// le labels, including +Inf), plus _sum and _count. +func Map(mp embed.MetricPoint) (MetricFamily, error) { + base := sanitizeName(mp.Name) + labels := sortedLabels(mp.Metadata.Attributes, targetLabels(mp.Metadata.Resource)) + tsMS := mp.Metadata.Timestamp.UnixMilli() + + fam := MetricFamily{Help: mp.Description, Unit: mp.Unit} + + switch mp.Type { + case embed.MetricTypeGauge, embed.MetricTypeSum: + fam.Type = TypeGauge + fam.Name = base + fam.Samples = []Sample{{Name: base, Labels: labels, Value: scalarValue(mp), TimestampMS: tsMS}} + case embed.MetricTypeCounter: + fam.Type = TypeCounter + fam.Name = withTotal(base) + fam.Samples = []Sample{{Name: fam.Name, Labels: labels, Value: scalarValue(mp), TimestampMS: tsMS}} + case embed.MetricTypeHistogram: + fam.Type = TypeHistogram + fam.Name = base + samples, err := histogramSamples(base, labels, mp, tsMS) + if err != nil { + return MetricFamily{}, err + } + fam.Samples = samples + default: + return MetricFamily{}, fmt.Errorf("prommap: unsupported metric type %q for %q", mp.Type, mp.Name) + } + return fam, nil +} + +// scalarValue returns the gauge/sum/counter value as a float64, preferring +// IntValue when set, then DoubleValue, else 0. +func scalarValue(mp embed.MetricPoint) float64 { + switch { + case mp.IntValue != nil: + return float64(*mp.IntValue) + case mp.DoubleValue != nil: + return *mp.DoubleValue + default: + return 0 + } +} + +// histogramSamples expands an OTel explicit-bucket histogram into cumulative +// _bucket series (le, plus +Inf), _sum, and _count. OTel counts are per-bucket +// with one overflow bucket, so len(counts) must be len(bounds)+1. +func histogramSamples(base string, labels []Label, mp embed.MetricPoint, tsMS int64) ([]Sample, error) { + bounds := mp.HistogramBucketBounds + counts := mp.HistogramBucketCounts + if len(counts) != len(bounds)+1 { + return nil, fmt.Errorf("prommap: histogram %q has %d bucket counts, want len(bounds)+1 = %d", mp.Name, len(counts), len(bounds)+1) + } + + samples := make([]Sample, 0, len(bounds)+3) + var cumulative uint64 + for i, bound := range bounds { + cumulative += counts[i] + samples = append(samples, Sample{ + Name: base + "_bucket", + Labels: withLE(labels, formatFloat(bound)), + Value: float64(cumulative), + TimestampMS: tsMS, + }) + } + // The overflow bucket becomes le="+Inf" and holds the full count. + cumulative += counts[len(counts)-1] + samples = append(samples, + Sample{Name: base + "_bucket", Labels: withLE(labels, "+Inf"), Value: float64(cumulative), TimestampMS: tsMS}, + Sample{Name: base + "_sum", Labels: labels, Value: mp.HistogramSum, TimestampMS: tsMS}, + // _count must equal the +Inf bucket, so derive it from the buckets + // rather than trusting HistogramCount to agree. + Sample{Name: base + "_count", Labels: labels, Value: float64(cumulative), TimestampMS: tsMS}, + ) + return samples, nil +} + +// sortedLabels converts metric attributes to sanitized, name-sorted labels. +// Target labels (job, instance) win over a same-named attribute. +func sortedLabels(attrs map[string]string, target []Label) []Label { + if len(attrs) == 0 && len(target) == 0 { + return nil + } + byName := make(map[string]string, len(attrs)+len(target)) + for k, v := range attrs { + byName[sanitizeLabelName(k)] = v + } + for _, l := range target { + byName[l.Name] = l.Value + } + return sortLabels(byName) +} + +func sortLabels(byName map[string]string) []Label { + labels := make([]Label, 0, len(byName)) + for k, v := range byName { + labels = append(labels, Label{Name: k, Value: v}) + } + sort.Slice(labels, func(i, j int) bool { return labels[i].Name < labels[j].Name }) + return labels +} + +// Resource keys the OTel-to-Prometheus compatibility spec promotes to job and +// instance. They feed the target labels and are left off target_info. +const ( + keyServiceName = "service.name" + keyServiceNamespace = "service.namespace" + keyServiceInstanceID = "service.instance.id" + keyHostName = "host.name" + keyTelemetrySource = "telemetry.source" +) + +// targetLabels derives job and instance from a resource. It follows the spec +// (job = [service.namespace/]service.name, instance = service.instance.id) +// and falls back to telemetry.source and host.name, which blitz generators +// set in place of service.*. Empty values are omitted. +func targetLabels(res map[string]any) []Label { + if len(res) == 0 { + return nil + } + job := resourceString(res, keyServiceName) + if ns := resourceString(res, keyServiceNamespace); job != "" && ns != "" { + job = ns + "/" + job + } + if job == "" { + job = resourceString(res, keyTelemetrySource) + } + instance := resourceString(res, keyServiceInstanceID) + if instance == "" { + instance = resourceString(res, keyHostName) + } + var labels []Label + if instance != "" { + labels = append(labels, Label{Name: "instance", Value: instance}) + } + if job != "" { + labels = append(labels, Label{Name: "job", Value: job}) + } + return labels +} + +// TargetInfo builds the target_info series for a point's resource: a gauge +// of 1 labeled with job, instance, and every other resource attribute. It +// reports false when the resource yields neither job nor instance. +func TargetInfo(mp embed.MetricPoint) (MetricFamily, bool) { + target := targetLabels(mp.Metadata.Resource) + if len(target) == 0 { + return MetricFamily{}, false + } + byName := make(map[string]string, len(mp.Metadata.Resource)+len(target)) + for k, v := range mp.Metadata.Resource { + switch k { + case keyServiceName, keyServiceNamespace, keyServiceInstanceID: + continue + } + byName[sanitizeLabelName(k)] = attrString(v) + } + for _, l := range target { + byName[l.Name] = l.Value + } + return MetricFamily{ + Name: "target_info", + Type: TypeGauge, + Help: "Target metadata", + Samples: []Sample{{ + Name: "target_info", + Labels: sortLabels(byName), + Value: 1, + TimestampMS: mp.Metadata.Timestamp.UnixMilli(), + }}, + }, true +} + +func resourceString(res map[string]any, key string) string { + v, ok := res[key] + if !ok { + return "" + } + return attrString(v) +} + +// attrString renders a resource value as the collector does (pcommon +// Value.AsString): strings as-is, scalars formatted, slices and maps as JSON. +func attrString(v any) string { + switch t := v.(type) { + case string: + return t + case nil: + return "" + case bool, int, int8, int16, int32, int64, uint, uint8, uint16, uint32, uint64, float32, float64: + return fmt.Sprint(t) + default: + b, err := json.Marshal(t) + if err != nil { + return fmt.Sprint(t) + } + return string(b) + } +} + +// withLE copies base and appends le last (Prometheus puts le last on a bucket). +func withLE(base []Label, le string) []Label { + out := make([]Label, len(base), len(base)+1) + copy(out, base) + return append(out, Label{Name: "le", Value: le}) +} + +// withTotal appends _total to a counter name unless it is already present. +func withTotal(name string) string { + if len(name) >= len("_total") && name[len(name)-len("_total"):] == "_total" { + return name + } + return name + "_total" +} + +// formatFloat renders a bucket boundary as the shortest round-trippable decimal. +func formatFloat(f float64) string { + return strconv.FormatFloat(f, 'g', -1, 64) +} + +// sanitizeName maps a name to a valid Prometheus metric name (colons allowed). +func sanitizeName(s string) string { return sanitize(s, true) } + +// sanitizeLabelName maps a key to a valid Prometheus label name (no colons). +func sanitizeLabelName(s string) string { return sanitize(s, false) } + +func sanitize(s string, allowColon bool) string { + if s == "" { + return s + } + out := make([]byte, 0, len(s)) + for i := 0; i < len(s); i++ { + c := s[i] + switch { + case c >= 'a' && c <= 'z', c >= 'A' && c <= 'Z', c == '_': + out = append(out, c) + case c == ':' && allowColon: + out = append(out, c) + case c >= '0' && c <= '9': + if i == 0 { + out = append(out, '_') + } + out = append(out, c) + default: + out = append(out, '_') + } + } + return string(out) +} diff --git a/internal/prommap/prommap_test.go b/internal/prommap/prommap_test.go new file mode 100644 index 00000000..5bd3a612 --- /dev/null +++ b/internal/prommap/prommap_test.go @@ -0,0 +1,186 @@ +package prommap + +import ( + "testing" + "time" + + "github.com/observiq/blitz/embed" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func i64(v int64) *int64 { return &v } +func f64(v float64) *float64 { return &v } +func ts() time.Time { return time.Unix(1700000000, 0) } + +const tsMS int64 = 1700000000000 + +// find returns the sample with the given series name, or fails. +func find(t *testing.T, fam MetricFamily, name string) Sample { + t.Helper() + for _, s := range fam.Samples { + if s.Name == name { + return s + } + } + t.Fatalf("no sample named %q in %+v", name, fam.Samples) + return Sample{} +} + +func TestMap_Gauge(t *testing.T) { + fam, err := Map(embed.MetricPoint{ + Name: "system.cpu.utilization", + Description: "CPU utilization", + Unit: "1", + Type: embed.MetricTypeGauge, + DoubleValue: f64(0.5), + Metadata: embed.MetricPointMetadata{ + Timestamp: ts(), + Attributes: map[string]string{"state": "idle", "cpu": "0"}, + }, + }) + require.NoError(t, err) + assert.Equal(t, "system_cpu_utilization", fam.Name) + assert.Equal(t, TypeGauge, fam.Type) + assert.Equal(t, "CPU utilization", fam.Help) + assert.Equal(t, "1", fam.Unit) + require.Len(t, fam.Samples, 1) + s := fam.Samples[0] + assert.Equal(t, "system_cpu_utilization", s.Name) + assert.Equal(t, 0.5, s.Value) + assert.Equal(t, tsMS, s.TimestampMS) + // Labels sorted by name. + assert.Equal(t, []Label{{"cpu", "0"}, {"state", "idle"}}, s.Labels) +} + +func TestMap_Counter_TotalSuffix(t *testing.T) { + fam, err := Map(embed.MetricPoint{ + Name: "http.requests", + Type: embed.MetricTypeCounter, + IntValue: i64(1027), + Metadata: embed.MetricPointMetadata{Timestamp: ts()}, + }) + require.NoError(t, err) + assert.Equal(t, TypeCounter, fam.Type) + assert.Equal(t, "http_requests_total", fam.Name) + require.Len(t, fam.Samples, 1) + assert.Equal(t, "http_requests_total", fam.Samples[0].Name) + assert.Equal(t, 1027.0, fam.Samples[0].Value) +} + +func TestMap_Counter_AlreadyTotal(t *testing.T) { + fam, err := Map(embed.MetricPoint{ + Name: "http_requests_total", + Type: embed.MetricTypeCounter, + IntValue: i64(1), + Metadata: embed.MetricPointMetadata{Timestamp: ts()}, + }) + require.NoError(t, err) + assert.Equal(t, "http_requests_total", fam.Name, "_total must not be doubled") +} + +func TestMap_Sum_IsGauge(t *testing.T) { + fam, err := Map(embed.MetricPoint{ + Name: "queue.size", + Type: embed.MetricTypeSum, + IntValue: i64(5), + Metadata: embed.MetricPointMetadata{Timestamp: ts()}, + }) + require.NoError(t, err) + assert.Equal(t, TypeGauge, fam.Type, "non-monotonic Sum maps to gauge") + assert.Equal(t, "queue_size", fam.Name, "gauge takes no _total suffix") + assert.Equal(t, 5.0, fam.Samples[0].Value) +} + +func TestMap_Histogram(t *testing.T) { + fam, err := Map(embed.MetricPoint{ + Name: "request.duration", + Type: embed.MetricTypeHistogram, + HistogramBucketBounds: []float64{0.1, 0.5, 1}, + HistogramBucketCounts: []uint64{1, 2, 3, 4}, // len == bounds+1 + HistogramSum: 3.3, + HistogramCount: 10, + Metadata: embed.MetricPointMetadata{Timestamp: ts(), Attributes: map[string]string{"route": "/x"}}, + }) + require.NoError(t, err) + assert.Equal(t, TypeHistogram, fam.Type) + assert.Equal(t, "request_duration", fam.Name) + + // Cumulative bucket counts, each with an le label plus the base route label. + b1 := find(t, fam, "request_duration_bucket") + _ = b1 // multiple _bucket samples share the name; assert via le below + + le := func(v string) Sample { + for _, s := range fam.Samples { + if s.Name != "request_duration_bucket" { + continue + } + for _, l := range s.Labels { + if l.Name == "le" && l.Value == v { + return s + } + } + } + t.Fatalf("no _bucket sample with le=%q", v) + return Sample{} + } + assert.Equal(t, 1.0, le("0.1").Value) + assert.Equal(t, 3.0, le("0.5").Value) // 1+2 + assert.Equal(t, 6.0, le("1").Value) // 1+2+3 + assert.Equal(t, 10.0, le("+Inf").Value) // 1+2+3+4 + + // le buckets carry the base labels too, sorted (le sorts after route). + assert.Equal(t, []Label{{"route", "/x"}, {"le", "0.1"}}, le("0.1").Labels) + + assert.Equal(t, 3.3, find(t, fam, "request_duration_sum").Value) + assert.Equal(t, 10.0, find(t, fam, "request_duration_count").Value) +} + +func TestMap_Histogram_BucketLengthMismatch(t *testing.T) { + _, err := Map(embed.MetricPoint{ + Name: "bad", + Type: embed.MetricTypeHistogram, + HistogramBucketBounds: []float64{0.1, 0.5}, + HistogramBucketCounts: []uint64{1, 2, 3, 4}, // should be len(bounds)+1 == 3 + Metadata: embed.MetricPointMetadata{Timestamp: ts()}, + }) + require.Error(t, err) +} + +func TestMap_NameAndLabelSanitization(t *testing.T) { + fam, err := Map(embed.MetricPoint{ + Name: "weird-name.with/chars", + Type: embed.MetricTypeGauge, + DoubleValue: f64(1), + Metadata: embed.MetricPointMetadata{ + Timestamp: ts(), + Attributes: map[string]string{"http.method": "GET", "1bad": "x"}, + }, + }) + require.NoError(t, err) + assert.Equal(t, "weird_name_with_chars", fam.Name) + got := map[string]string{} + for _, l := range fam.Samples[0].Labels { + got[l.Name] = l.Value + } + assert.Equal(t, "GET", got["http_method"], "dotted label name sanitized") + assert.Equal(t, "x", got["_1bad"], "label name starting with a digit gets a leading underscore") +} + +func TestMap_Deterministic(t *testing.T) { + mp := embed.MetricPoint{ + Name: "m", + Type: embed.MetricTypeGauge, + DoubleValue: f64(2), + Metadata: embed.MetricPointMetadata{ + Timestamp: ts(), + Attributes: map[string]string{"b": "2", "a": "1", "c": "3"}, + }, + } + a, err := Map(mp) + require.NoError(t, err) + b, err := Map(mp) + require.NoError(t, err) + assert.Equal(t, a, b) + assert.Equal(t, []Label{{"a", "1"}, {"b", "2"}, {"c", "3"}}, a.Samples[0].Labels) +} diff --git a/internal/prommap/target_test.go b/internal/prommap/target_test.go new file mode 100644 index 00000000..073ef47a --- /dev/null +++ b/internal/prommap/target_test.go @@ -0,0 +1,137 @@ +package prommap + +import ( + "testing" + + "github.com/observiq/blitz/embed" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func gaugeWithResource(res map[string]any) embed.MetricPoint { + return embed.MetricPoint{ + Name: "system.memory.usage", + Type: embed.MetricTypeGauge, + IntValue: i64(42), + Metadata: embed.MetricPointMetadata{ + Timestamp: ts(), + Attributes: map[string]string{"state": "used"}, + Resource: res, + }, + } +} + +// Distinct hosts must produce distinct series: host.name becomes instance and +// telemetry.source becomes job. +func TestMap_ResourceBecomesJobAndInstance(t *testing.T) { + fam, err := Map(gaugeWithResource(map[string]any{ + "host.name": "athena", + "telemetry.source": "hostmetrics", + })) + require.NoError(t, err) + assert.Equal(t, []Label{ + {"instance", "athena"}, + {"job", "hostmetrics"}, + {"state", "used"}, + }, fam.Samples[0].Labels) +} + +func TestMap_DistinctHostsDistinctSeries(t *testing.T) { + a, err := Map(gaugeWithResource(map[string]any{"host.name": "athena", "telemetry.source": "hostmetrics"})) + require.NoError(t, err) + b, err := Map(gaugeWithResource(map[string]any{"host.name": "hermes", "telemetry.source": "hostmetrics"})) + require.NoError(t, err) + assert.NotEqual(t, a.Samples[0].Labels, b.Samples[0].Labels) +} + +// service.* takes precedence per the OTel-to-Prometheus compatibility spec. +func TestMap_ServiceAttributesTakePrecedence(t *testing.T) { + fam, err := Map(gaugeWithResource(map[string]any{ + "service.namespace": "shop", + "service.name": "cart", + "service.instance.id": "cart-7", + "host.name": "athena", + "telemetry.source": "hostmetrics", + })) + require.NoError(t, err) + assert.Equal(t, []Label{ + {"instance", "cart-7"}, + {"job", "shop/cart"}, + {"state", "used"}, + }, fam.Samples[0].Labels) +} + +// No resource means no job/instance labels (existing behavior preserved). +func TestMap_NoResourceNoTargetLabels(t *testing.T) { + fam, err := Map(gaugeWithResource(nil)) + require.NoError(t, err) + assert.Equal(t, []Label{{"state", "used"}}, fam.Samples[0].Labels) +} + +func TestTargetInfo(t *testing.T) { + fam, ok := TargetInfo(gaugeWithResource(map[string]any{ + "host.name": "athena", + "telemetry.source": "hostmetrics", + "os.type": "linux", + "host.ip": []string{"10.0.0.1", "10.0.0.2"}, + })) + require.True(t, ok) + assert.Equal(t, "target_info", fam.Name) + assert.Equal(t, TypeGauge, fam.Type) + require.Len(t, fam.Samples, 1) + s := fam.Samples[0] + assert.Equal(t, "target_info", s.Name) + assert.Equal(t, 1.0, s.Value) + assert.Equal(t, tsMS, s.TimestampMS) + assert.Equal(t, []Label{ + {"host_ip", `["10.0.0.1","10.0.0.2"]`}, + {"host_name", "athena"}, + {"instance", "athena"}, + {"job", "hostmetrics"}, + {"os_type", "linux"}, + {"telemetry_source", "hostmetrics"}, + }, s.Labels) +} + +// service.* keys feed job/instance and are excluded from target_info's labels. +func TestTargetInfo_ExcludesServiceKeys(t *testing.T) { + fam, ok := TargetInfo(gaugeWithResource(map[string]any{ + "service.name": "cart", + "service.instance.id": "cart-7", + "os.type": "linux", + })) + require.True(t, ok) + assert.Equal(t, []Label{ + {"instance", "cart-7"}, + {"job", "cart"}, + {"os_type", "linux"}, + }, fam.Samples[0].Labels) +} + +func TestTargetInfo_NoResource(t *testing.T) { + _, ok := TargetInfo(gaugeWithResource(nil)) + assert.False(t, ok) +} + +// _count must equal the le="+Inf" bucket even when HistogramCount disagrees +// with the bucket counts (Prometheus invariant). +func TestMap_Histogram_CountMatchesInf(t *testing.T) { + fam, err := Map(embed.MetricPoint{ + Name: "latency", + Type: embed.MetricTypeHistogram, + HistogramBucketBounds: []float64{1, 5}, + HistogramBucketCounts: []uint64{2, 3, 4}, + HistogramCount: 99, + HistogramSum: 20, + Metadata: embed.MetricPointMetadata{Timestamp: ts()}, + }) + require.NoError(t, err) + var inf float64 + for _, s := range fam.Samples { + if s.Name == "latency_bucket" && s.Labels[len(s.Labels)-1].Value == "+Inf" { + inf = s.Value + } + } + assert.Equal(t, 9.0, inf) + assert.Equal(t, 9.0, find(t, fam, "latency_count").Value) +} diff --git a/package/completions/blitz.bash b/package/completions/blitz.bash index b4ceb56e..bb20afb2 100644 --- a/package/completions/blitz.bash +++ b/package/completions/blitz.bash @@ -407,8 +407,6 @@ _blitz_help() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -693,8 +691,6 @@ _blitz_library_diff() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -980,8 +976,6 @@ _blitz_library_extract() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -1265,8 +1259,6 @@ _blitz_library_ls() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -1550,8 +1542,6 @@ _blitz_library_path() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -1835,8 +1825,6 @@ _blitz_library_search() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -2120,8 +2108,6 @@ _blitz_library_show() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -2411,8 +2397,6 @@ _blitz_library() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -2696,8 +2680,6 @@ _blitz_version() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=") @@ -2984,8 +2966,6 @@ _blitz_root_command() two_word_flags+=("--generator-hostmetrics-rate") flags+=("--generator-hostmetrics-scrapers=") two_word_flags+=("--generator-hostmetrics-scrapers") - flags+=("--generator-hostmetrics-workers=") - two_word_flags+=("--generator-hostmetrics-workers") flags+=("--generator-json-rate=") two_word_flags+=("--generator-json-rate") flags+=("--generator-json-type=")