Datasource/grafana-pyroscope: Add rate aggregation for cumulative profiles (#108546)

This commit is contained in:
Marc Sanmiquel
2025-08-18 17:08:34 +02:00
committed by GitHub
parent c508b01cc5
commit 9b800d63b1
6 changed files with 584 additions and 24 deletions
@@ -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
@@ -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"
}
}
@@ -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
}
}
@@ -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
})
}
+108 -13
View File
@@ -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{
@@ -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 {