diff --git a/pkg/tsdb/grafana-pyroscope-datasource/instance.go b/pkg/tsdb/grafana-pyroscope-datasource/instance.go index 978b7148bda..7943707b39a 100644 --- a/pkg/tsdb/grafana-pyroscope-datasource/instance.go +++ b/pkg/tsdb/grafana-pyroscope-datasource/instance.go @@ -79,6 +79,9 @@ func (d *PyroscopeDatasource) CallResource(ctx context.Context, req *backend.Cal if req.Path == "labelValues" { return d.labelValues(ctx, req, sender) } + if req.Path == "profileMetadata" { + return d.profileMetadata(ctx, req, sender) + } return sender.Send(&backend.CallResourceResponse{ Status: 404, }) @@ -216,6 +219,32 @@ func (d *PyroscopeDatasource) labelValues(ctx context.Context, req *backend.Call return nil } +// profileMetadata returns the embedded profile-metrics.json data containing metadata +// for all known profile types, including their aggregation type (cumulative/instant), +// units, descriptions, and grouping information. +func (d *PyroscopeDatasource) profileMetadata(ctx context.Context, _ *backend.CallResourceRequest, sender backend.CallResourceResponseSender) error { + ctxLogger := logger.FromContext(ctx) + + registry := GetProfileMetadataRegistry() + + jsonData, err := json.Marshal(registry.profiles) + if err != nil { + ctxLogger.Error("Failed to marshal profile metadata", "error", err, "function", logEntrypoint()) + return sender.Send(&backend.CallResourceResponse{ + Status: 500, + Body: []byte(`{"error": "Failed to marshal profile metadata"}`), + }) + } + + return sender.Send(&backend.CallResourceResponse{ + Status: 200, + Body: jsonData, + Headers: map[string][]string{ + "Content-Type": {"application/json"}, + }, + }) +} + // QueryData handles multiple queries and returns multiple responses. // req contains the queries []DataQuery (where each query contains RefID as a unique identifier). // The QueryDataResponse contains a map of RefID to the response for each query, and each response diff --git a/pkg/tsdb/grafana-pyroscope-datasource/profile-metrics.json b/pkg/tsdb/grafana-pyroscope-datasource/profile-metrics.json new file mode 100644 index 00000000000..5912e37175e --- /dev/null +++ b/pkg/tsdb/grafana-pyroscope-datasource/profile-metrics.json @@ -0,0 +1,162 @@ +{ + "block:contentions:count:contentions:count": { + "id": "block:contentions:count:contentions:count", + "description": "Number of blocking contentions", + "type": "contentions", + "group": "block", + "unit": "short", + "aggregationType": "cumulative" + }, + "block:delay:nanoseconds:contentions:count": { + "id": "block:delay:nanoseconds:contentions:count", + "description": "Time spent in blocking delays", + "type": "delay", + "group": "block", + "unit": "ns", + "aggregationType": "cumulative" + }, + "goroutine:goroutine:count:goroutine:count": { + "id": "goroutine:goroutine:count:goroutine:count", + "description": "Number of goroutines", + "type": "goroutine", + "group": "goroutine", + "unit": "short", + "aggregationType": "instant" + }, + "goroutines:goroutine:count:goroutine:count": { + "id": "goroutines:goroutine:count:goroutine:count", + "description": "Number of goroutines", + "type": "goroutine", + "group": "goroutine", + "unit": "short", + "aggregationType": "instant" + }, + "memory:alloc_in_new_tlab_bytes:bytes::": { + "id": "memory:alloc_in_new_tlab_bytes:bytes::", + "description": "Size of memory allocated inside Thread-Local Allocation Buffers (TLAB)", + "type": "alloc_in_new_tlab_bytes", + "group": "memory", + "unit": "bytes", + "aggregationType": "cumulative" + }, + "memory:alloc_in_new_tlab_objects:count::": { + "id": "memory:alloc_in_new_tlab_objects:count::", + "description": "Number of objects allocated inside Thread-Local Allocation Buffers (TLAB)", + "type": "alloc_in_new_tlab_objects", + "group": "memory", + "unit": "short", + "aggregationType": "cumulative" + }, + "memory:alloc_objects:count:space:bytes": { + "id": "memory:alloc_objects:count:space:bytes", + "description": "Number of objects allocated", + "type": "alloc_objects", + "group": "memory", + "unit": "short", + "aggregationType": "cumulative" + }, + "memory:alloc_space:bytes:space:bytes": { + "id": "memory:alloc_space:bytes:space:bytes", + "description": "Size of memory allocated in the heap", + "type": "alloc_space", + "group": "memory", + "unit": "bytes", + "aggregationType": "cumulative" + }, + "memory:inuse_objects:count:space:bytes": { + "id": "memory:inuse_objects:count:space:bytes", + "description": "Number of objects currently in use", + "type": "inuse_objects", + "group": "memory", + "unit": "short", + "aggregationType": "instant" + }, + "memory:inuse_space:bytes:space:bytes": { + "id": "memory:inuse_space:bytes:space:bytes", + "description": "Size of memory currently in use", + "type": "inuse_space", + "group": "memory", + "unit": "bytes", + "aggregationType": "instant" + }, + "mutex:contentions:count:contentions:count": { + "id": "mutex:contentions:count:contentions:count", + "description": "Number of observed mutex contentions", + "type": "contentions", + "group": "mutex", + "unit": "short", + "aggregationType": "cumulative" + }, + "mutex:delay:nanoseconds:contentions:count": { + "id": "mutex:delay:nanoseconds:contentions:count", + "description": "Time spent waiting due to mutex contentions", + "type": "delay", + "group": "mutex", + "unit": "ns", + "aggregationType": "cumulative" + }, + "process_cpu:alloc_samples:count:cpu:nanoseconds": { + "id": "process_cpu:alloc_samples:count:cpu:nanoseconds", + "description": "Number of memory allocation samples during CPU time", + "type": "alloc_samples", + "group": "memory", + "unit": "short", + "aggregationType": "cumulative" + }, + "process_cpu:alloc_size:bytes:cpu:nanoseconds": { + "id": "process_cpu:alloc_size:bytes:cpu:nanoseconds", + "description": "Size of memory allocated during CPU time", + "type": "alloc_size", + "group": "alloc_size", + "unit": "bytes", + "aggregationType": "cumulative" + }, + "process_cpu:cpu:nanoseconds:cpu:nanoseconds": { + "id": "process_cpu:cpu:nanoseconds:cpu:nanoseconds", + "description": "CPU time consumed", + "type": "cpu", + "group": "process_cpu", + "unit": "ns", + "aggregationType": "cumulative" + }, + "process_cpu:exception:count:cpu:nanoseconds": { + "id": "process_cpu:exception:count:cpu:nanoseconds", + "description": "Number of exceptions within the sampled CPU time", + "type": "exceptions", + "group": "exceptions", + "unit": "short", + "aggregationType": "cumulative" + }, + "process_cpu:lock_count:count:cpu:nanoseconds": { + "id": "process_cpu:lock_count:count:cpu:nanoseconds", + "description": "Number of lock acquisitions attempted during CPU time", + "type": "lock_count", + "group": "locks", + "unit": "short", + "aggregationType": "instant" + }, + "process_cpu:lock_time:nanoseconds:cpu:nanoseconds": { + "id": "process_cpu:lock_time:nanoseconds:cpu:nanoseconds", + "description": "Cumulative time spent acquiring locks", + "type": "lock_time", + "group": "locks", + "unit": "ns", + "aggregationType": "cumulative" + }, + "process_cpu:samples:count::milliseconds": { + "id": "process_cpu:samples:count::milliseconds", + "description": "Number of process samples collected", + "type": "samples", + "group": "process_cpu", + "unit": "short", + "aggregationType": "cumulative" + }, + "process_cpu:samples:count:cpu:nanoseconds": { + "id": "process_cpu:samples:count:cpu:nanoseconds", + "description": "Number of samples collected over CPU time", + "type": "samples", + "group": "process_cpu", + "unit": "short", + "aggregationType": "instant" + } +} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/profile_metadata.go b/pkg/tsdb/grafana-pyroscope-datasource/profile_metadata.go new file mode 100644 index 00000000000..6f778acef49 --- /dev/null +++ b/pkg/tsdb/grafana-pyroscope-datasource/profile_metadata.go @@ -0,0 +1,89 @@ +package pyroscope + +import ( + _ "embed" + "encoding/json" + "sync" + + "github.com/grafana/grafana-plugin-sdk-go/backend" + "github.com/grafana/grafana-plugin-sdk-go/backend/log" +) + +//go:embed profile-metrics.json +var profileMetricsJSON []byte + +type ProfileMetadata struct { + ID string `json:"id"` + Description string `json:"description"` + Type string `json:"type"` + Group string `json:"group"` + Unit string `json:"unit"` + AggregationType string `json:"aggregationType"` +} + +type ProfileMetadataRegistry struct { + profiles map[string]*ProfileMetadata + mu sync.RWMutex + logger log.Logger +} + +var ( + registry *ProfileMetadataRegistry + registryOnce sync.Once +) + +// GetProfileMetadataRegistry returns the singleton instance of the profile metadata registry +func GetProfileMetadataRegistry() *ProfileMetadataRegistry { + registryOnce.Do(func() { + registry = &ProfileMetadataRegistry{ + profiles: make(map[string]*ProfileMetadata), + logger: backend.NewLoggerWith("logger", "tsdb.pyroscope.profile-metadata"), + } + registry.loadProfileMetadata() + }) + return registry +} + +// loadProfileMetadata loads the profile metadata from the embedded JSON file +func (r *ProfileMetadataRegistry) loadProfileMetadata() { + var profilesMap map[string]*ProfileMetadata + err := json.Unmarshal(profileMetricsJSON, &profilesMap) + if err != nil { + r.logger.Error("Failed to parse embedded profile-metrics.json", "error", err) + return + } + + r.mu.Lock() + defer r.mu.Unlock() + r.profiles = profilesMap + r.logger.Info("Loaded profile metadata", "count", len(profilesMap)) +} + +// GetProfileMetadata returns the metadata for a given profile type ID +func (r *ProfileMetadataRegistry) GetProfileMetadata(profileTypeID string) *ProfileMetadata { + r.mu.RLock() + defer r.mu.RUnlock() + return r.profiles[profileTypeID] +} + +// IsCumulativeProfile returns true if the profile type requires rate calculation +func (r *ProfileMetadataRegistry) IsCumulativeProfile(profileTypeID string) bool { + metadata := r.GetProfileMetadata(profileTypeID) + if metadata == nil { + r.logger.Debug("Profile metadata not found, using fallback logic", "profileTypeID", profileTypeID) + return isCumulativeProfileUnitFallback(getUnits(profileTypeID)) + } + return metadata.AggregationType == "cumulative" +} + +// isCumulativeProfileUnitFallback is the (old) fallback logic for unknown profile types +func isCumulativeProfileUnitFallback(unit string) bool { + switch unit { + case "ns": + return true + case "bytes": + return true + default: + return false + } +} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/profile_metadata_test.go b/pkg/tsdb/grafana-pyroscope-datasource/profile_metadata_test.go new file mode 100644 index 00000000000..8c43884969f --- /dev/null +++ b/pkg/tsdb/grafana-pyroscope-datasource/profile_metadata_test.go @@ -0,0 +1,49 @@ +package pyroscope + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestProfileMetadataRegistry(t *testing.T) { + registry := GetProfileMetadataRegistry() + + t.Run("CPU profile is cumulative", func(t *testing.T) { + result := registry.IsCumulativeProfile("process_cpu:cpu:nanoseconds:cpu:nanoseconds") + require.True(t, result) + }) + + t.Run("Memory allocation is cumulative", func(t *testing.T) { + result := registry.IsCumulativeProfile("memory:alloc_space:bytes:space:bytes") + require.True(t, result) + }) + + t.Run("Goroutines are instant", func(t *testing.T) { + result := registry.IsCumulativeProfile("goroutine:goroutine:count:goroutine:count") + require.False(t, result) + }) + + t.Run("Memory in-use is instant", func(t *testing.T) { + result := registry.IsCumulativeProfile("memory:inuse_space:bytes:space:bytes") + require.False(t, result) + }) + + t.Run("Edge case: mutex contentions count is cumulative", func(t *testing.T) { + result := registry.IsCumulativeProfile("mutex:contentions:count:contentions:count") + require.True(t, result) + }) + + t.Run("Edge case: memory alloc objects count is cumulative", func(t *testing.T) { + result := registry.IsCumulativeProfile("memory:alloc_objects:count:space:bytes") + require.True(t, result) + }) + + t.Run("Unknown profile falls back to unit-based logic", func(t *testing.T) { + result := registry.IsCumulativeProfile("unknown:profile:nanoseconds:test:test") + require.True(t, result) // ns should be treated as cumulative by fallback + + result = registry.IsCumulativeProfile("unknown:profile:count:test:test") + require.False(t, result) // count/short should be treated as instant by fallback + }) +} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/query.go b/pkg/tsdb/grafana-pyroscope-datasource/query.go index c2243c7af05..506f9de33ee 100644 --- a/pkg/tsdb/grafana-pyroscope-datasource/query.go +++ b/pkg/tsdb/grafana-pyroscope-datasource/query.go @@ -76,7 +76,6 @@ func (d *PyroscopeDatasource) query(ctx context.Context, pCtx backend.PluginCont logger.Error("Failed to parse the MinStep using default", "MinStep", dsJson.MinStep, "function", logEntrypoint()) } } - logger.Debug("Sending SelectSeriesRequest", "queryModel", qm, "function", logEntrypoint()) seriesResp, err := d.client.GetSeries( gCtx, profileTypeId, @@ -96,7 +95,8 @@ func (d *PyroscopeDatasource) query(ctx context.Context, pCtx backend.PluginCont // add the frames to the response. responseMutex.Lock() withAnnotations := qm.Annotations != nil && *qm.Annotations - frames, err := seriesToDataFrames(seriesResp, withAnnotations) + stepDuration := math.Max(query.Interval.Seconds(), parsedInterval.Seconds()) + frames, err := seriesToDataFrames(seriesResp, withAnnotations, stepDuration, profileTypeId) if err != nil { span.RecordError(err) span.SetStatus(codes.Error, err.Error()) @@ -136,7 +136,24 @@ func (d *PyroscopeDatasource) query(ctx context.Context, pCtx backend.PluginCont var frame *data.Frame if profileResp != nil { - frame = responseToDataFrames(profileResp) + var dsJson dsJsonModel + err := json.Unmarshal(pCtx.DataSourceInstanceSettings.JSONData, &dsJson) + if err != nil { + span.RecordError(err) + span.SetStatus(codes.Error, err.Error()) + return fmt.Errorf("error unmarshaling datasource json model: %v", err) + } + + parsedInterval := time.Second * 15 + if dsJson.MinStep != "" { + parsedInterval, err = gtime.ParseDuration(dsJson.MinStep) + if err != nil { + parsedInterval = time.Second * 15 + logger.Error("Failed to parse the MinStep using default", "MinStep", dsJson.MinStep, "function", logEntrypoint()) + } + } + stepDuration := math.Max(query.Interval.Seconds(), parsedInterval.Seconds()) + frame = responseToDataFrames(profileResp, stepDuration, profileTypeId) // If query called with streaming on then return a channel // to subscribe on a client-side and consume updates from a plugin. @@ -173,9 +190,9 @@ func (d *PyroscopeDatasource) query(ctx context.Context, pCtx backend.PluginCont // responseToDataFrames turns Pyroscope response to data.Frame. We encode the data into a nested set format where we have // [level, value, label] columns and by ordering the items in a depth first traversal order we can recreate the whole // tree back. -func responseToDataFrames(resp *ProfileResponse) *data.Frame { +func responseToDataFrames(resp *ProfileResponse, stepDurationSec float64, profileTypeID string) *data.Frame { tree := levelsToTree(resp.Flamebearer.Levels, resp.Flamebearer.Names) - return treeToNestedSetDataFrame(tree, resp.Units) + return treeToNestedSetDataFrame(tree, resp.Units, stepDurationSec, profileTypeID) } // START_OFFSET is offset of the bar relative to previous sibling @@ -337,9 +354,17 @@ type CustomMeta struct { // where ordering the items in depth first order and knowing the level/depth of each item we can recreate the // parent - child relationship without explicitly needing parent/child column, and we can later just iterate over the // dataFrame to again basically walking depth first over the tree/profile. -func treeToNestedSetDataFrame(tree *ProfileTree, unit string) *data.Frame { +func treeToNestedSetDataFrame(tree *ProfileTree, unit string, stepDurationSec float64, profileTypeID string) *data.Frame { frame := data.NewFrame("response") - frame.Meta = &data.FrameMeta{PreferredVisualization: "flamegraph"} + frameMeta := &data.FrameMeta{PreferredVisualization: "flamegraph"} + + // Add metadata when rate calculation is applied + if isCumulativeProfile(profileTypeID) && stepDurationSec > 0 { + frameMeta.Custom = map[string]interface{}{ + "rateCalculated": true, + } + } + frame.Meta = frameMeta levelField := data.NewField("level", nil, []int64{}) valueField := data.NewField("value", nil, []int64{}) @@ -356,8 +381,17 @@ func treeToNestedSetDataFrame(tree *ProfileTree, unit string) *data.Frame { if tree != nil { walkTree(tree, func(tree *ProfileTree) { levelField.Append(int64(tree.Level)) - valueField.Append(tree.Value) - selfField.Append(tree.Self) + + // Apply rate calculation for cumulative profiles + value := tree.Value + self := tree.Self + if isCumulativeProfile(profileTypeID) && stepDurationSec > 0 { + value = int64(float64(value) / stepDurationSec) + self = int64(float64(self) / stepDurationSec) + } + + valueField.Append(value) + selfField.Append(self) labelField.Append(tree.Name) }) } @@ -433,14 +467,53 @@ func (ta *TimedAnnotation) getValue() string { return ta.Annotation.Value } -func seriesToDataFrames(resp *SeriesResponse, withAnnotations bool) ([]*data.Frame, error) { +// isCumulativeProfile determines if a profile type requires rate calculation using the metadata registry +func isCumulativeProfile(profileTypeID string) bool { + registry := GetProfileMetadataRegistry() + return registry.IsCumulativeProfile(profileTypeID) +} + +// isCPUTimeProfile determines if a profile type represents CPU time in nanoseconds +func isCPUTimeProfile(profileTypeID string) bool { + registry := GetProfileMetadataRegistry() + metadata := registry.GetProfileMetadata(profileTypeID) + if metadata != nil { + // Check if it's CPU time (unit is nanoseconds and type contains cpu) + return metadata.Unit == "ns" && (metadata.Type == "cpu" || metadata.Group == "process_cpu") + } + return false +} + +// convertToRateUnit converts profile units to appropriate rate units when rate calculation is applied +func convertToRateUnit(originalUnit string) string { + switch originalUnit { + case "bytes": + return "binBps" + case "short": + return "ops" + case "ns": + return "ns" + default: + return originalUnit + } +} + +func seriesToDataFrames(resp *SeriesResponse, withAnnotations bool, stepDurationSec float64, profileTypeID string) ([]*data.Frame, error) { frames := make([]*data.Frame, 0, len(resp.Series)) annotations := make([]*TimedAnnotation, 0) for _, series := range resp.Series { // We create separate data frames as the series may not have the same length frame := data.NewFrame("series") - frame.Meta = &data.FrameMeta{PreferredVisualization: "graph"} + frameMeta := &data.FrameMeta{PreferredVisualization: "graph"} + + // Add metadata when rate calculation is applied + if isCumulativeProfile(profileTypeID) && stepDurationSec > 0 { + frameMeta.Custom = map[string]interface{}{ + "rateCalculated": true, + } + } + frame.Meta = frameMeta fields := make(data.Fields, 0, 2) timeField := data.NewField("time", nil, []time.Time{}) @@ -451,13 +524,35 @@ func seriesToDataFrames(resp *SeriesResponse, withAnnotations bool) ([]*data.Fra labels[label.Name] = label.Value } + // Determine display unit - convert units for rate-calculated cumulative profiles + displayUnit := resp.Units + if isCumulativeProfile(profileTypeID) && stepDurationSec > 0 { + if isCPUTimeProfile(profileTypeID) { + displayUnit = "cores" + } else { + // Convert other cumulative profile units to rate units + displayUnit = convertToRateUnit(resp.Units) + } + } + valueField := data.NewField(resp.Label, labels, []float64{}) - valueField.Config = &data.FieldConfig{Unit: resp.Units} + valueField.Config = &data.FieldConfig{Unit: displayUnit} fields = append(fields, valueField) for _, point := range series.Points { timeField.Append(time.UnixMilli(point.Timestamp)) - valueField.Append(point.Value) + + // Apply rate calculation for cumulative profiles + value := point.Value + if isCumulativeProfile(profileTypeID) && stepDurationSec > 0 { + value = value / stepDurationSec + + // Convert CPU nanoseconds to cores + if isCPUTimeProfile(profileTypeID) { + value = value / 1e9 + } + } + valueField.Append(value) if withAnnotations { for _, a := range point.Annotations { annotations = append(annotations, &TimedAnnotation{ diff --git a/pkg/tsdb/grafana-pyroscope-datasource/query_test.go b/pkg/tsdb/grafana-pyroscope-datasource/query_test.go index def783b917f..67addaaedd4 100644 --- a/pkg/tsdb/grafana-pyroscope-datasource/query_test.go +++ b/pkg/tsdb/grafana-pyroscope-datasource/query_test.go @@ -5,10 +5,11 @@ import ( "testing" "time" + "github.com/stretchr/testify/require" + "github.com/grafana/grafana-plugin-sdk-go/backend" "github.com/grafana/grafana-plugin-sdk-go/data" typesv1 "github.com/grafana/pyroscope/api/gen/proto/go/types/v1" - "github.com/stretchr/testify/require" ) // This is where the tests for the datasource backend live. @@ -130,7 +131,7 @@ func Test_profileToDataFrame(t *testing.T) { }, Units: "short", } - frame := responseToDataFrames(profile) + frame := responseToDataFrames(profile, 15.0, "goroutine:goroutine:count:goroutine:count") require.Equal(t, 4, len(frame.Fields)) require.Equal(t, data.NewField("level", nil, []int64{0, 1, 1}), frame.Fields[0]) require.Equal(t, data.NewField("value", nil, []int64{20, 10, 5}).SetConfig(&data.FieldConfig{Unit: "short"}), frame.Fields[1]) @@ -202,7 +203,7 @@ func Test_treeToNestedDataFrame(t *testing.T) { }, } - frame := treeToNestedSetDataFrame(tree, "short") + frame := treeToNestedSetDataFrame(tree, "short", 15.0, "goroutine:goroutine:count:goroutine:count") labelConfig := &data.FieldConfig{ TypeConfig: &data.FieldTypeConfig{ @@ -221,10 +222,52 @@ func Test_treeToNestedDataFrame(t *testing.T) { }) t.Run("nil profile tree", func(t *testing.T) { - frame := treeToNestedSetDataFrame(nil, "short") + frame := treeToNestedSetDataFrame(nil, "short", 15.0, "goroutine:goroutine:count:goroutine:count") require.Equal(t, 4, len(frame.Fields)) require.Equal(t, 0, frame.Fields[0].Len()) }) + + t.Run("rateCalculated metadata for cumulative profile", func(t *testing.T) { + tree := &ProfileTree{ + Value: 100, Level: 0, Self: 1, Name: "root", + } + frame := treeToNestedSetDataFrame(tree, "short", 15.0, "process_cpu:cpu:nanoseconds:cpu:nanoseconds") + require.NotNil(t, frame.Meta) + require.NotNil(t, frame.Meta.Custom) + custom := frame.Meta.Custom.(map[string]interface{}) + require.Equal(t, true, custom["rateCalculated"]) + }) + + t.Run("no rateCalculated metadata for instant profile", func(t *testing.T) { + tree := &ProfileTree{ + Value: 100, Level: 0, Self: 1, Name: "root", + } + frame := treeToNestedSetDataFrame(tree, "short", 15.0, "goroutine:goroutine:count:goroutine:count") + require.NotNil(t, frame.Meta) + require.Nil(t, frame.Meta.Custom) + }) + + t.Run("CPU time keeps original units for tree data", func(t *testing.T) { + tree := &ProfileTree{ + Value: 3000000000, Level: 0, Self: 1500000000, Name: "root", // 3s total, 1.5s self in nanoseconds + } + // Test CPU profile (should keep nanoseconds for flamegraph, no unit conversion) + frame := treeToNestedSetDataFrame(tree, "ns", 15.0, "process_cpu:cpu:nanoseconds:cpu:nanoseconds") + + // Check unit remains as nanoseconds (no conversion for flamegraphs) + require.Equal(t, "ns", frame.Fields[1].Config.Unit) + require.Equal(t, "ns", frame.Fields[2].Config.Unit) + + // Check values were rate calculated but not unit converted: 3000000000/15 = 200000000, 1500000000/15 = 100000000 + require.Equal(t, int64(200000000), frame.Fields[1].At(0)) + require.Equal(t, int64(100000000), frame.Fields[2].At(0)) + + // Check metadata shows rate was calculated + require.NotNil(t, frame.Meta) + require.NotNil(t, frame.Meta.Custom) + custom := frame.Meta.Custom.(map[string]interface{}) + require.Equal(t, true, custom["rateCalculated"]) + }) } func Test_seriesToDataFrameAnnotations(t *testing.T) { @@ -253,7 +296,7 @@ func Test_seriesToDataFrameAnnotations(t *testing.T) { Label: "samples", } - frames, err := seriesToDataFrames(series, true) + frames, err := seriesToDataFrames(series, true, 15.0, "goroutine:goroutine:count:goroutine:count") require.NoError(t, err) require.Equal(t, 1, len(frames)) require.Equal(t, 2, len(frames[0].Fields)) @@ -278,7 +321,7 @@ func Test_seriesToDataFrameAnnotations(t *testing.T) { }, } - frames, err := seriesToDataFrames(series, false) + frames, err := seriesToDataFrames(series, false, 15.0, "goroutine:goroutine:count:goroutine:count") require.NoError(t, err) require.Equal(t, 1, len(frames)) }) @@ -302,7 +345,7 @@ func Test_seriesToDataFrameAnnotations(t *testing.T) { }, } - frames, err := seriesToDataFrames(series, true) + frames, err := seriesToDataFrames(series, true, 15.0, "goroutine:goroutine:count:goroutine:count") require.NoError(t, err) require.Equal(t, 2, len(frames)) @@ -348,7 +391,7 @@ func Test_seriesToDataFrameAnnotations(t *testing.T) { }, } - frames, err := seriesToDataFrames(series, true) + frames, err := seriesToDataFrames(series, true, 15.0, "goroutine:goroutine:count:goroutine:count") require.NoError(t, err) require.Equal(t, 2, len(frames)) @@ -375,7 +418,7 @@ func Test_seriesToDataFrame(t *testing.T) { Units: "short", Label: "samples", } - frames, err := seriesToDataFrames(series, true) + frames, err := seriesToDataFrames(series, true, 15.0, "goroutine:goroutine:count:goroutine:count") require.NoError(t, err) require.Equal(t, 2, len(frames[0].Fields)) require.Equal(t, data.NewField("time", nil, []time.Time{time.UnixMilli(1000), time.UnixMilli(2000)}), frames[0].Fields[0]) @@ -390,7 +433,7 @@ func Test_seriesToDataFrame(t *testing.T) { Label: "samples", } - frames, err = seriesToDataFrames(series, true) + frames, err = seriesToDataFrames(series, true, 15.0, "goroutine:goroutine:count:goroutine:count") require.NoError(t, err) require.Equal(t, data.NewField("samples", map[string]string{"app": "bar"}, []float64{30, 10}).SetConfig(&data.FieldConfig{Unit: "short"}), frames[0].Fields[1]) }) @@ -404,7 +447,7 @@ func Test_seriesToDataFrame(t *testing.T) { Units: "short", Label: "samples", } - frames, err := seriesToDataFrames(resp, true) + frames, err := seriesToDataFrames(resp, true, 15.0, "goroutine:goroutine:count:goroutine:count") require.NoError(t, err) require.Equal(t, 2, len(frames)) require.Equal(t, 2, len(frames[0].Fields)) @@ -412,6 +455,99 @@ func Test_seriesToDataFrame(t *testing.T) { require.Equal(t, data.NewField("samples", map[string]string{"foo": "bar"}, []float64{30, 10}).SetConfig(&data.FieldConfig{Unit: "short"}), frames[0].Fields[1]) require.Equal(t, data.NewField("samples", map[string]string{"foo": "baz"}, []float64{30, 10}).SetConfig(&data.FieldConfig{Unit: "short"}), frames[1].Fields[1]) }) + + t.Run("rateCalculated metadata for cumulative profile", func(t *testing.T) { + series := &SeriesResponse{ + Series: []*Series{ + {Labels: []*LabelPair{}, Points: []*Point{{Timestamp: int64(1000), Value: 30}, {Timestamp: int64(2000), Value: 10}}}, + }, + Units: "ns", + Label: "cpu", + } + frames, err := seriesToDataFrames(series, false, 15.0, "process_cpu:cpu:nanoseconds:cpu:nanoseconds") + require.NoError(t, err) + require.Equal(t, 1, len(frames)) + require.NotNil(t, frames[0].Meta) + require.NotNil(t, frames[0].Meta.Custom) + custom := frames[0].Meta.Custom.(map[string]interface{}) + require.Equal(t, true, custom["rateCalculated"]) + }) + + t.Run("no rateCalculated metadata for instant profile", func(t *testing.T) { + series := &SeriesResponse{ + Series: []*Series{ + {Labels: []*LabelPair{}, Points: []*Point{{Timestamp: int64(1000), Value: 30}, {Timestamp: int64(2000), Value: 10}}}, + }, + Units: "short", + Label: "goroutines", + } + // Test instant profile (should not have rateCalculated metadata) + frames, err := seriesToDataFrames(series, false, 15.0, "goroutine:goroutine:count:goroutine:count") + require.NoError(t, err) + require.Equal(t, 1, len(frames)) + require.NotNil(t, frames[0].Meta) + require.Nil(t, frames[0].Meta.Custom) + }) + + t.Run("CPU time conversion to cores", func(t *testing.T) { + series := &SeriesResponse{ + Series: []*Series{ + {Labels: []*LabelPair{}, Points: []*Point{{Timestamp: int64(1000), Value: 3000000000}, {Timestamp: int64(2000), Value: 1500000000}}}, // 3s and 1.5s in nanoseconds + }, + Units: "ns", + Label: "cpu", + } + // should convert nanoseconds to cores and set unit to "cores" + frames, err := seriesToDataFrames(series, false, 15.0, "process_cpu:cpu:nanoseconds:cpu:nanoseconds") + require.NoError(t, err) + require.Equal(t, 1, len(frames)) + + require.Equal(t, "cores", frames[0].Fields[1].Config.Unit) + + // Check values were converted: 3000000000/15/1e9 = 0.2 cores/sec, 1500000000/15/1e9 = 0.1 cores/sec + values := fieldValues[float64](frames[0].Fields[1]) + require.Equal(t, []float64{0.2, 0.1}, values) + }) + + t.Run("Memory allocation unit conversion to bytes/sec", func(t *testing.T) { + series := &SeriesResponse{ + Series: []*Series{ + {Labels: []*LabelPair{}, Points: []*Point{{Timestamp: int64(1000), Value: 150000000}, {Timestamp: int64(2000), Value: 300000000}}}, // 150 MB, 300 MB + }, + Units: "bytes", + Label: "memory_alloc", + } + // should convert bytes to binBps and apply rate calculation + frames, err := seriesToDataFrames(series, false, 15.0, "memory:alloc_space:bytes:space:bytes") + require.NoError(t, err) + require.Equal(t, 1, len(frames)) + + require.Equal(t, "binBps", frames[0].Fields[1].Config.Unit) + + // Check values were rate calculated: 150000000/15 = 10000000, 300000000/15 = 20000000 + values := fieldValues[float64](frames[0].Fields[1]) + require.Equal(t, []float64{10000000, 20000000}, values) + }) + + t.Run("Count-based profile unit conversion to ops/sec", func(t *testing.T) { + series := &SeriesResponse{ + Series: []*Series{ + {Labels: []*LabelPair{}, Points: []*Point{{Timestamp: int64(1000), Value: 1500}, {Timestamp: int64(2000), Value: 3000}}}, // 1500, 3000 contentions + }, + Units: "short", + Label: "contentions", + } + // should convert short to ops and apply rate calculation + frames, err := seriesToDataFrames(series, false, 15.0, "mutex:contentions:count:contentions:count") + require.NoError(t, err) + require.Equal(t, 1, len(frames)) + + require.Equal(t, "ops", frames[0].Fields[1].Config.Unit) + + // Check values were rate calculated: 1500/15 = 100, 3000/15 = 200 + values := fieldValues[float64](frames[0].Fields[1]) + require.Equal(t, []float64{100, 200}, values) + }) } type FakeClient struct {