diff --git a/pkg/metric/BUILD b/pkg/metric/BUILD index f7e35c12d..04235af55 100644 --- a/pkg/metric/BUILD +++ b/pkg/metric/BUILD @@ -25,6 +25,7 @@ go_library( "//pkg/log", "//pkg/prometheus", "//pkg/sync", + "@org_golang_google_protobuf//encoding/protojson:go_default_library", "@org_golang_google_protobuf//types/known/timestamppb", ], ) diff --git a/pkg/metric/metric_test.go b/pkg/metric/metric_test.go index af8022829..0bbdf1ca8 100644 --- a/pkg/metric/metric_test.go +++ b/pkg/metric/metric_test.go @@ -1140,29 +1140,24 @@ func TestMetricProfiling(t *testing.T) { t.Fatalf("failed to open file a second time: %v", err) } - // Check the header - h := adler32.New() - lines := bufio.NewScanner(f) - expectedHeader := "Time (ns)\t" + strings.Join(test.metricNames, "\t") - if test.lossy { - expectedHeader = MetricsPrefix + expectedHeader - } - var header string - for header == "" && lines.Scan() { - header = lines.Text() - } - if header != expectedHeader { - t.Fatalf("header mismatch: got '%s' want '%s'", header, expectedHeader) - } - h.Write([]byte(strings.TrimPrefix(header, MetricsPrefix) + "\n")) - - // Check that data looks sane: + // Check that the log looks sane: + // - Header should match what we expect. + // - We should have one metadata line per metric. + // - We should have one start time line. + // - If in lossy mode, we should have a hash line. + // - If we have a hash, it should match the one computed by this test. // - Each timestamp should always be bigger. // - 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") prevTS := uint64(0) prevValues := make([]uint64, numMetrics) numDatapoints := 0 var hashLine string + gotMetadataFor := make(map[string]struct{}, numMetrics) + gotHeader := false + gotStartTime := false for lines.Scan() { line := lines.Text() if line == "" { @@ -1182,6 +1177,34 @@ func TestMetricProfiling(t *testing.T) { continue } h.Write([]byte(line + "\n")) + if strings.HasPrefix(line, MetricsMetaIndicator) { + line = strings.TrimPrefix(line, MetricsMetaIndicator) + components := strings.Split(line, "\t") + if len(components) != 2 { + t.Fatalf("got %d components in metadata line %q, want 2", len(components), line) + } + // We only verify that the metadata is present, not its contents. + if components[0] == "" || components[1] == "" { + t.Fatalf("got empty metadata line: %q", line) + } + gotMetadataFor[components[0]] = struct{}{} + continue + } + if strings.HasPrefix(line, MetricsStartTimeIndicator) { + if gotStartTime { + t.Fatalf("got multiple start time lines") + } + gotStartTime = true + continue + } + if !gotHeader { + // This line must be the header. + if line != expectedHeader { + t.Fatalf("got header %q, want %q", line, expectedHeader) + } + gotHeader = true + continue + } numDatapoints++ items := strings.Split(line, "\t") if len(items) != (numMetrics + 1) { @@ -1220,6 +1243,20 @@ func TestMetricProfiling(t *testing.T) { t.Errorf("incorrect final metric value: got %d, want %d", prevValues[i], expected) } } + if len(gotMetadataFor) != numMetrics { + t.Errorf("got metadata for %d metrics, want %d", len(gotMetadataFor), numMetrics) + } + for _, metricName := range test.metricNames { + if _, ok := gotMetadataFor[metricName]; !ok { + t.Errorf("did not get metadata for metric %q", metricName) + } + } + if !gotStartTime { + t.Error("did not get start time metadata") + } + if !gotHeader { + t.Error("did not get header") + } if test.lossy { if hashLine == "" { t.Fatal("lossy writer output does not have expected hash line") diff --git a/pkg/metric/profiling_metric.go b/pkg/metric/profiling_metric.go index 24de96529..5b717a60e 100644 --- a/pkg/metric/profiling_metric.go +++ b/pkg/metric/profiling_metric.go @@ -24,6 +24,7 @@ import ( "strings" "time" + "google.golang.org/protobuf/encoding/protojson" "gvisor.dev/gvisor/pkg/atomicbitops" "gvisor.dev/gvisor/pkg/log" "gvisor.dev/gvisor/pkg/prometheus" @@ -37,6 +38,14 @@ const ( // MetricsHashIndicator is prepended before the hash of the metrics // data at the end of the metrics stream. MetricsHashIndicator = "ADLER32\t" + // TimeColumn is the column header for the time column. + TimeColumn = "Time (ns)" + // MetricsMetaIndicator is prepended before every metrics metadata line + // after metricsPrefix. + MetricsMetaIndicator = "META\t" + // MetricsStartTimeIndicator is prepended before the start time of the + // metrics collection. + MetricsStartTimeIndicator = "START_TIME\t" ) var ( @@ -58,6 +67,9 @@ var ( // before it's written to the writer. type snapshots struct { numMetrics int + // startTime is the time at which collection started, as reported by + // CheapNowNano() which does *not* start counting from the epoch time. + startTime int64 // ringbuffer is used to store metric data. ringbuffer [][]uint64 // curWriterIndex is the ringbuffer index currently being read by the @@ -113,8 +125,9 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption } var values []func(fieldValues ...*FieldValue) uint64 - var header strings.Builder - header.WriteString("Time (ns)") + var headers []string + var columnHeaders strings.Builder + columnHeaders.WriteString(TimeColumn) numMetrics := 0 if len(opts.Metrics) > 0 { @@ -131,8 +144,18 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption // 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) } - header.WriteRune('\t') - header.WriteString(name) + var metricMetadataHeader strings.Builder + metricMetadataHeader.WriteString(MetricsMetaIndicator) + metricMetadataHeader.WriteString(name) + metricMetadataHeader.WriteRune('\t') + metricMetadata, err := protojson.MarshalOptions{Multiline: false}.Marshal(m.metadata) + if err != nil { + return fmt.Errorf("failed to marshal metric schema for metric %q: %w", name, err) + } + metricMetadataHeader.Write(metricMetadata) + headers = append(headers, metricMetadataHeader.String()) + columnHeaders.WriteRune('\t') + columnHeaders.WriteString(name) values = append(values, m.value) } } else { @@ -141,6 +164,11 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption } return fmt.Errorf("a value for --profiling-metrics was not specified; also no conditionally compiled metrics found, consider compiling runsc with --go_tag=condmetric_profiling") } + headers = append( + headers, + fmt.Sprintf("%s%d", MetricsStartTimeIndicator, time.Now().UnixNano()), + columnHeaders.String(), + ) if !profilingMetricsStarted.CompareAndSwap(false, true) { return errors.New("profiling metrics have already been started") @@ -160,14 +188,14 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption stopProfilingMetrics = make(chan bool, 1) doneProfilingMetrics = make(chan bool, 1) writeCh := make(chan writeReq, snapshotRingbufferSize) - + s.startTime = CheapNowNano() go collectProfilingMetrics(&s, values, opts.Rate, writeCh) if opts.Lossy { lossySink := newLossyBufferedWriter(opts.Sink) - go writeProfilingMetrics[*lossyBufferedWriter[T]](lossySink, &s, header.String(), writeCh) + go writeProfilingMetrics[*lossyBufferedWriter[T]](lossySink, &s, headers, writeCh) } else { bufferedSink := newBufferedWriter(opts.Sink) - go writeProfilingMetrics[*bufferedWriter[T]](bufferedSink, &s, header.String(), writeCh) + go writeProfilingMetrics[*bufferedWriter[T]](bufferedSink, &s, headers, writeCh) } log.Infof("Profiling metrics started.") @@ -182,7 +210,6 @@ func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldVa numEntries := s.numMetrics + 1 // to account for the timestamp ringbufferIdx := 0 curSnapshot := 0 - startTime := CheapNowNano() // getNewRingbufferIdx will block until the writer indicates that some part // of the ringbuffer is available for writing. getNewRingbufferIdx := func() { @@ -216,7 +243,7 @@ func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldVa } collectStart := CheapNowNano() - timestamp := time.Duration(collectStart - startTime) + timestamp := time.Duration(collectStart - s.startTime) base := curSnapshot * numEntries s.ringbuffer[ringbufferIdx][base] = uint64(timestamp) for i := 1; i < numEntries; i++ { @@ -397,10 +424,12 @@ 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, header string, writeReqs <-chan writeReq) { +func writeProfilingMetrics[T bufferedMetricsWriter](sink T, s *snapshots, headers []string, writeReqs <-chan writeReq) { numEntries := s.numMetrics + 1 - sink.WriteString(header) - sink.NewLine() + for _, header := range headers { + sink.WriteString(header) + sink.NewLine() + } for req := range writeReqs { s.curWriterIndex.Store(int32(req.ringbufferIdx)) for i := 0; i < req.numLines; i++ { @@ -416,6 +445,7 @@ func writeProfilingMetrics[T bufferedMetricsWriter](sink T, s *snapshots, header } } sink.Close() + doneProfilingMetrics <- true close(doneProfilingMetrics) profilingMetricsStarted.Store(false)