diff --git a/pkg/metric/metric_test.go b/pkg/metric/metric_test.go index 06d777906..1dd11296c 100644 --- a/pkg/metric/metric_test.go +++ b/pkg/metric/metric_test.go @@ -1062,10 +1062,16 @@ func TestFieldMapperMustUseSameValuePointer(t *testing.T) { } func TestMetricProfiling(t *testing.T) { + fieldVal1 := &FieldValue{"val1"} + fieldVal2 := &FieldValue{"val2"} + fieldVal3 := &FieldValue{"val3"} + for _, test := range []struct { name string profilingMetricsFlag string metricNames []string + firstMetricFields []Field + expectedHeader string // If unspecified, computed automatically for non-field metrics. lossy bool incrementMetricBy []uint64 numIterations uint64 @@ -1104,6 +1110,30 @@ func TestMetricProfiling(t *testing.T) { numIterations: 100, errOnStartProfiling: false, }, + { + name: "single metric with 1 field", + profilingMetricsFlag: "/foo", + metricNames: []string{"/foo"}, + firstMetricFields: []Field{ + NewField("field1", fieldVal1, fieldVal2, fieldVal3), + }, + expectedHeader: TimeColumn + "\t/foo[val1]\t/foo[val2]\t/foo[val3]", + incrementMetricBy: []uint64{1337}, + numIterations: 100, + errOnStartProfiling: false, + }, + { + name: "multiple metrics and one has fields", + profilingMetricsFlag: "/foo,/bar", + metricNames: []string{"/foo", "/bar"}, + firstMetricFields: []Field{ + NewField("field1", fieldVal1, fieldVal2, fieldVal3), + }, + expectedHeader: TimeColumn + "\t/foo[val1]\t/foo[val2]\t/foo[val3]\t/bar", + incrementMetricBy: []uint64{42, 27}, + numIterations: 100, + errOnStartProfiling: false, + }, { name: "mismatched names", profilingMetricsFlag: "/foo,/fighter,/big/test/baz,/metric", @@ -1116,20 +1146,42 @@ func TestMetricProfiling(t *testing.T) { t.Run(test.name, func(t *testing.T) { defer resetTest() const profilingRate = 1000 * time.Microsecond - numMetrics := len(test.metricNames) + numEntries := 0 + for i := range test.metricNames { + if i == 0 { + fm, err := newFieldMapper(test.firstMetricFields...) + if err != nil { + t.Fatalf("newFieldMapper err: got %v wanted nil", err) + } + numEntries += fm.numKeys() + } else { + numEntries++ + } + } - metrics := make([]*Uint64Metric, numMetrics) + metrics := make([]*Uint64Metric, len(test.metricNames)) for i, m := range test.metricNames { + var fields []Field + if i == 0 { + for _, f := range test.firstMetricFields { + fields = append(fields, f) + } + } newMetric, err := NewUint64Metric(m, Uint64Metadata{ Cumulative: true, Sync: true, Description: fooDescription, + Fields: fields, }) metrics[i] = newMetric if err != nil { t.Fatalf("NewUint64Metric got err '%v' want nil", err) } } + firstMetricFieldMapper, err := newFieldMapper(test.firstMetricFields...) + if err != nil { + t.Fatalf("newFieldMapper err: got %v wanted nil", err) + } f, err := os.CreateTemp(t.TempDir(), "profiling-metrics") if err != nil { @@ -1158,7 +1210,15 @@ func TestMetricProfiling(t *testing.T) { // Generate some test data for i := 0; i < int(test.numIterations); i++ { for metricIdx, m := range metrics { - m.IncrementBy(test.incrementMetricBy[metricIdx]) + if metricIdx == 0 { + for fieldKey := 0; fieldKey < firstMetricFieldMapper.numKeys(); fieldKey++ { + fieldValues := make([]*FieldValue, len(test.firstMetricFields)) + firstMetricFieldMapper.keyToMultiFieldInPlace(fieldKey, fieldValues) + m.IncrementBy(test.incrementMetricBy[metricIdx], fieldValues...) + } + } else { + m.IncrementBy(test.incrementMetricBy[metricIdx]) + } } // Give time to the collector to record it. time.Sleep(profilingRate) @@ -1181,15 +1241,18 @@ func TestMetricProfiling(t *testing.T) { // - Each metric value should be at least as big as the previous. h := adler32.New() lines := bufio.NewScanner(f) - expectedHeader := TimeColumn + "\t" + strings.Join(test.metricNames, "\t") + expectedHeader := test.expectedHeader + if expectedHeader == "" { + expectedHeader = TimeColumn + "\t" + strings.Join(test.metricNames, "\t") + } if test.lossy { expectedHeader += "\tChecksum" } prevTS := uint64(0) - prevValues := make([]uint64, numMetrics) + prevValues := make([]uint64, numEntries) numDatapoints := 0 var hashLine string - gotMetadataFor := make(map[string]struct{}, numMetrics) + gotMetadataFor := make(map[string]struct{}, len(test.metricNames)) gotHeader := false gotStartTime := false gotStats := false @@ -1264,8 +1327,8 @@ func TestMetricProfiling(t *testing.T) { } numDatapoints++ items := strings.Split(line, "\t") - if len(items) != (numMetrics + 1) { - t.Fatalf("incorrect number of items on line '%s': got %d, want %d", line, len(items), numMetrics+1) + if len(items) != (numEntries + 1) { + t.Fatalf("incorrect number of items on line '%s': got %d, want %d", line, len(items), numEntries+1) } // Check timestamp ts, err := strconv.ParseUint(items[0], 10, 64) @@ -1278,7 +1341,7 @@ func TestMetricProfiling(t *testing.T) { prevTS = ts // Check metric values - for i := 1; i <= numMetrics; i++ { + for i := 1; i <= numEntries; i++ { m, err := strconv.ParseUint(items[i], 10, 64) if err != nil { t.Errorf("m ParseUint error on line '%s': got '%v' want nil", line, err) @@ -1294,14 +1357,22 @@ func TestMetricProfiling(t *testing.T) { t.Errorf("numDatapoints: got %d, want at least %d", numDatapoints, expectedMinNumDatapoints) } // Check that the final total for each metric is correct - for i := 0; i < numMetrics; i++ { - expected := test.numIterations * test.incrementMetricBy[i] - if prevValues[i] != expected { - t.Errorf("incorrect final metric value: got %d, want %d", prevValues[i], expected) + for metricIdx := range metrics { + expected := test.numIterations * test.incrementMetricBy[metricIdx] + if metricIdx == 0 { + for fieldKey := 0; fieldKey < firstMetricFieldMapper.numKeys(); fieldKey++ { + if prevValues[fieldKey] != expected { + t.Errorf("incorrect final metric value: got %d, want %d", prevValues[fieldKey], expected) + } + } + } else { + if itemIdx := firstMetricFieldMapper.numKeys() + metricIdx - 1; prevValues[itemIdx] != expected { + t.Errorf("incorrect final metric value: got %d, want %d", prevValues[itemIdx], expected) + } } } - if len(gotMetadataFor) != numMetrics { - t.Errorf("got metadata for %d metrics, want %d", len(gotMetadataFor), numMetrics) + if len(gotMetadataFor) != len(test.metricNames) { + t.Errorf("got metadata for %d metrics, want %d", len(gotMetadataFor), len(test.metricNames)) } for _, metricName := range test.metricNames { if _, ok := gotMetadataFor[metricName]; !ok { diff --git a/pkg/metric/profiling_metric.go b/pkg/metric/profiling_metric.go index 3d52703c6..b178066a9 100644 --- a/pkg/metric/profiling_metric.go +++ b/pkg/metric/profiling_metric.go @@ -37,8 +37,7 @@ import ( const ( // snapshotBufferSize is the number of snapshots within one item of the // ringbuffer. Increasing this number means less context-switching - // overhead between collector and writer goroutines, but worse time - // precision, as the precise time is refreshed every this many snapshots. + // overhead between collector and writer goroutines, but more memory usage. snapshotBufferSize = 1024 // snapshotRingbufferSize is the number of items in the ringbuffer. // Increasing this number means the writer has more slack to catch up @@ -72,7 +71,7 @@ type CollectionStats struct { // CollectionRate is the rate at which the metrics are meant to be // collected. - CollectionRateNanos uint64 `json:"collection_rate"` + CollectionRateNanos uint64 `json:"collection_rate_nanos"` // CheapStartNanos is the time at which the collector started in nanoseconds, // as returned by CheapNowNano. @@ -89,7 +88,7 @@ type CollectionStats struct { // between when the collector was meant to start collecting a metric snapshot // vs when it actually collected it. It should be divided by TotalSnapshots // to get the average sleep timing error. - TotalSleepTimingErrorNanos uint64 `json:"total_sleep_timing_error"` + TotalSleepTimingErrorNanos uint64 `json:"total_sleep_timing_error_nanos"` // TotalCollectionTimingError is the running sum of time spent doing actual // metric collection, i.e. retrieving numerical values from metrics. @@ -100,7 +99,7 @@ type CollectionStats struct { // not actually being the case. // It should be divided by TotalSnapshots to get the average // per-collection-cycle collection timing error. - TotalCollectionTimingErrorNanos uint64 `json:"total_collection_timing_error"` + TotalCollectionTimingErrorNanos uint64 `json:"total_collection_timing_error_nanos"` // NumBackoffSleeps is the number of times the collector had to back off // because the writer was too slow. @@ -108,7 +107,7 @@ type CollectionStats struct { // TotalBackoffSleep is the running sum of time the collector had to back // off because the writer was too slow. - TotalBackoffSleepNanos uint64 `json:"total_backoff_sleep"` + TotalBackoffSleepNanos uint64 `json:"total_backoff_sleep_nanos"` } var ( @@ -130,7 +129,7 @@ var ( // snapshots is used to as temporary storage of metric data // before it's written to the writer. type snapshots struct { - numMetrics int + numEntries int // startTime is the time at which collection started in nanoseconds. startTime int64 // ringbuffer is used to store metric data. @@ -180,6 +179,11 @@ type ProfilingMetricsOptions[T ProfilingMetricsWriter] struct { Rate time.Duration } +// valueFunc returns the value of a single metric with a single field value +// combination. For distributions, it returns a single statistic of that +// distribution. +type valueFunc func() uint64 + // StartProfilingMetrics checks the ProfilingMetrics runsc flags and creates // goroutines responsible for outputting the profiling metric data. // @@ -193,26 +197,19 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption return errors.New("metric initialization is not complete") } - var values []func(fieldValues ...*FieldValue) uint64 + var valueFuncs []valueFunc var headers []string var columnHeaders strings.Builder columnHeaders.WriteString(TimeColumn) - numMetrics := 0 + numEntries := 0 if len(opts.Metrics) > 0 { - metrics := strings.Split(opts.Metrics, ",") - numMetrics = len(metrics) - - for _, name := range metrics { + for _, name := range strings.Split(opts.Metrics, ",") { name := strings.TrimSpace(name) m, ok := allMetrics.uint64Metrics[name] if !ok { return fmt.Errorf("given profiling metric name '%s' does not correspond to a registered Uint64 metric", name) } - if len(m.fields) > 0 { - // TODO(b/240280155): Add support for field values. - return fmt.Errorf("will not profile metric '%s' because it has metric fields which are not supported", name) - } var metricMetadataHeader strings.Builder metricMetadataHeader.WriteString(MetricsMetaIndicator) metricMetadataHeader.WriteString(name) @@ -223,9 +220,31 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption } metricMetadataHeader.Write(metricMetadata) headers = append(headers, metricMetadataHeader.String()) - columnHeaders.WriteRune('\t') - columnHeaders.WriteString(name) - values = append(values, m.value) + fieldMapper, err := newFieldMapper(m.fields...) + if err != nil { + return fmt.Errorf("failed to create field mapper for metric %q: %w", name, err) + } + metricValueFunc := m.value + for fieldKey := 0; fieldKey < fieldMapper.numKeys(); fieldKey++ { + fieldValues := make([]*FieldValue, len(m.fields)) + fieldMapper.keyToMultiFieldInPlace(fieldKey, fieldValues) + columnHeaders.WriteRune('\t') + columnHeaders.WriteString(name) + if len(fieldValues) > 0 { + columnHeaders.WriteRune('[') + for i, fieldValue := range fieldValues { + if i > 0 { + columnHeaders.WriteRune(',') + } + columnHeaders.WriteString(fieldValue.Value) + } + columnHeaders.WriteRune(']') + } + valueFuncs = append(valueFuncs, func() uint64 { + return metricValueFunc(fieldValues...) + }) + numEntries++ + } } if opts.Lossy { columnHeaders.WriteString("\tChecksum") @@ -246,7 +265,7 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption return errors.New("profiling metrics have already been started") } s := snapshots{ - numMetrics: numMetrics, + numEntries: numEntries, ringbuffer: make([][]uint64, snapshotRingbufferSize), // curWriterIndex is initialized to a valid index so that the // collector cannot use up all indices before the writer even has @@ -254,7 +273,7 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption curWriterIndex: atomicbitops.FromInt32(snapshotRingbufferSize - 1), } for i := 0; i < snapshotRingbufferSize; i++ { - s.ringbuffer[i] = make([]uint64, snapshotBufferSize*(numMetrics+1)) + s.ringbuffer[i] = make([]uint64, snapshotBufferSize*(numEntries+1)) } // Truncate the underlying sink if possible to delete any past profiling @@ -273,7 +292,7 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption CollectionRateNanos: uint64(opts.Rate.Nanoseconds()), CheapStartNanos: uint64(cheapStartTime), } - go collectProfilingMetrics(&s, values, cheapStartTime, opts.Rate, writeCh, &stats) + go collectProfilingMetrics(&s, valueFuncs, cheapStartTime, opts.Rate, writeCh, &stats) if opts.Lossy { lossySink := newLossyBufferedWriter(opts.Sink) go writeProfilingMetrics[*lossyBufferedWriter[T]](lossySink, &s, headers, writeCh, &stats) @@ -288,13 +307,13 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption // collectProfilingMetrics will send metrics to the writeCh until it receives a // signal via the stopProfilingMetrics channel. -func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldValue) uint64, cheapStartTime int64, profilingRate time.Duration, writeCh chan<- writeReq, stats *CollectionStats) { +func collectProfilingMetrics(s *snapshots, valueFuncs []valueFunc, cheapStartTime int64, profilingRate time.Duration, writeCh chan<- writeReq, stats *CollectionStats) { defer close(writeCh) stats.mu.Lock() defer stats.mu.Unlock() - numEntries := s.numMetrics + 1 // to account for the timestamp + numEntries := s.numEntries + 1 // to account for the timestamp ringbufferIdx := 0 curSnapshot := 0 var beforeCollectionTimestamp int64 @@ -365,7 +384,7 @@ func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldVa ringBuf := s.ringbuffer[ringbufferIdx] base := curSnapshot * numEntries for i := 1; i < numEntries; i++ { - ringBuf[base+i] = values[i-1]() + ringBuf[base+i] = valueFuncs[i-1]() } afterCollectionTimestamp := CheapNowNano() middleCollectionTimestamp := (beforeCollectionTimestamp + afterCollectionTimestamp) / 2 @@ -598,7 +617,7 @@ func (w *lossyBufferedWriter[T]) Close() error { // writeProfilingMetrics will write to the ProfilingMetricsWriter on every // request via writeReqs, until writeReqs is closed. func writeProfilingMetrics[T bufferedMetricsWriter](sink T, s *snapshots, headers []string, writeReqs <-chan writeReq, stats *CollectionStats) { - numEntries := s.numMetrics + 1 + numEntries := s.numEntries + 1 for _, header := range headers { sink.WriteString(header) sink.NewLine()