Pyroscope: Annotation support for series queries (#104130)

* Pyroscope: Add annotations frame to series response

* Adapt to API change, add tests

* Run make lint-go

* Fix conflicts after rebase

* Add annotation via a separate data frame

* Process annotations fully at the datasource

* Add mod owner for go-humanize

* Pyroscope: Annotations in Query Response can be optional

---------

Co-authored-by: Piotr Jamróz <pm.jamroz@gmail.com>
This commit is contained in:
Aleksandar Petrov
2025-05-28 10:42:19 +02:00
committed by GitHub
co-authored by Piotr Jamróz
parent ea0e49a6e6
commit 0b8252fd7c
11 changed files with 535 additions and 17 deletions
@@ -0,0 +1,133 @@
package pyroscope
import (
"encoding/json"
"fmt"
"time"
"github.com/dustin/go-humanize"
"github.com/grafana/grafana-plugin-sdk-go/data"
)
// profileAnnotationKey represents the key for different types of annotations
type profileAnnotationKey string
const (
// profileAnnotationKeyThrottled is the key for throttling annotations
profileAnnotationKeyThrottled profileAnnotationKey = "pyroscope.ingest.throttled"
)
// ProfileAnnotation represents the parsed annotation data
type ProfileAnnotation struct {
Body ProfileThrottledAnnotation `json:"body"`
}
// ProfileThrottledAnnotation contains throttling information
type ProfileThrottledAnnotation struct {
PeriodType string `json:"periodType"`
PeriodLimitMb float64 `json:"periodLimitMb"`
LimitResetTime int64 `json:"limitResetTime"`
SamplingPeriodSec float64 `json:"samplingPeriodSec"`
SamplingRequests int64 `json:"samplingRequests"`
UsageGroup string `json:"usageGroup"`
}
// processedProfileAnnotation represents a processed annotation ready for display
type processedProfileAnnotation struct {
text string
time int64
timeEnd int64
isRegion bool
duplicateTracker int64
}
// grafanaAnnotationData holds slices of processed annotation data
type grafanaAnnotationData struct {
times []time.Time
timeEnds []time.Time
texts []string
isRegions []bool
}
// convertAnnotation converts a Pyroscope profile annotation into a Grafana annotation
func convertAnnotation(timedAnnotation *TimedAnnotation, duplicateTracker int64) (*processedProfileAnnotation, error) {
if timedAnnotation.getKey() != string(profileAnnotationKeyThrottled) {
// Currently we only support throttling annotations
return nil, nil
}
var profileAnnotation ProfileAnnotation
err := json.Unmarshal([]byte(timedAnnotation.getValue()), &profileAnnotation)
if err != nil {
return nil, fmt.Errorf("error parsing annotation data: %w", err)
}
throttlingInfo := profileAnnotation.Body
if duplicateTracker == throttlingInfo.LimitResetTime {
return nil, nil
}
limit := humanize.IBytes(uint64(throttlingInfo.PeriodLimitMb * 1024 * 1024))
return &processedProfileAnnotation{
text: fmt.Sprintf("Ingestion limit (%s/%s) reached", limit, throttlingInfo.PeriodType),
time: timedAnnotation.Timestamp,
timeEnd: throttlingInfo.LimitResetTime * 1000,
isRegion: throttlingInfo.LimitResetTime < time.Now().Unix(),
duplicateTracker: throttlingInfo.LimitResetTime,
}, nil
}
// processAnnotations processes a slice of TimedAnnotation and returns grafanaAnnotationData
func processAnnotations(timedAnnotations []*TimedAnnotation) (*grafanaAnnotationData, error) {
result := &grafanaAnnotationData{
times: []time.Time{},
timeEnds: []time.Time{},
texts: []string{},
isRegions: []bool{},
}
var duplicateTracker int64
for _, timedAnnotation := range timedAnnotations {
if timedAnnotation == nil || timedAnnotation.Annotation == nil {
continue
}
processed, err := convertAnnotation(timedAnnotation, duplicateTracker)
if err != nil {
return nil, err
}
if processed != nil {
result.times = append(result.times, time.UnixMilli(processed.time))
result.timeEnds = append(result.timeEnds, time.UnixMilli(processed.timeEnd))
result.isRegions = append(result.isRegions, processed.isRegion)
result.texts = append(result.texts, processed.text)
duplicateTracker = processed.duplicateTracker
}
}
return result, nil
}
// createAnnotationFrame creates a data frame for annotations
func createAnnotationFrame(annotations []*TimedAnnotation) (*data.Frame, error) {
annotationData, err := processAnnotations(annotations)
if err != nil {
return nil, err
}
timeField := data.NewField("time", nil, annotationData.times)
timeEndField := data.NewField("timeEnd", nil, annotationData.timeEnds)
textField := data.NewField("text", nil, annotationData.texts)
isRegionField := data.NewField("isRegion", nil, annotationData.isRegions)
colorField := data.NewField("color", nil, make([]string, len(annotationData.times)))
frame := data.NewFrame("annotations")
frame.Fields = data.Fields{timeField, timeEndField, textField, isRegionField, colorField}
frame.SetMeta(&data.FrameMeta{
DataTopic: data.DataTopicAnnotations,
})
return frame, nil
}
@@ -0,0 +1,188 @@
package pyroscope
import (
"testing"
"time"
"github.com/grafana/grafana-plugin-sdk-go/data"
typesv1 "github.com/grafana/pyroscope/api/gen/proto/go/types/v1"
"github.com/stretchr/testify/require"
)
func TestConvertAnnotation(t *testing.T) {
rawAnnotation := `{"body":{"periodType":"day","periodLimitMb":1024,"limitResetTime":1609459200}}`
t.Run("processes valid annotation", func(t *testing.T) {
timedAnnotation := &TimedAnnotation{
Timestamp: 1609455600000,
Annotation: &typesv1.ProfileAnnotation{
Key: string(profileAnnotationKeyThrottled),
Value: rawAnnotation,
},
}
processed, err := convertAnnotation(timedAnnotation, 0)
require.NoError(t, err)
require.NotNil(t, processed)
require.Contains(t, processed.text, "Ingestion limit (1.0 GiB/day) reached")
require.Contains(t, processed.text, "day")
require.Equal(t, int64(1609455600000), processed.time)
require.Equal(t, int64(1609459200000), processed.timeEnd) // LimitResetTime * 1000
require.Equal(t, int64(1609459200), processed.duplicateTracker)
})
t.Run("ignores non-throttling annotations", func(t *testing.T) {
timedAnnotation := &TimedAnnotation{
Timestamp: 1000,
Annotation: &typesv1.ProfileAnnotation{
Key: "some.other.key",
Value: `{"test":"value"}`,
},
}
processed, err := convertAnnotation(timedAnnotation, 0)
require.NoError(t, err)
require.Nil(t, processed)
})
t.Run("handles invalid annotation data", func(t *testing.T) {
timedAnnotation := &TimedAnnotation{
Timestamp: 1000,
Annotation: &typesv1.ProfileAnnotation{
Key: string(profileAnnotationKeyThrottled),
Value: `invalid json`,
},
}
processed, err := convertAnnotation(timedAnnotation, 0)
require.Error(t, err)
require.Nil(t, processed)
require.Contains(t, err.Error(), "error parsing annotation data")
})
t.Run("skips duplicate annotations", func(t *testing.T) {
timedAnnotation := &TimedAnnotation{
Timestamp: 1000,
Annotation: &typesv1.ProfileAnnotation{
Key: string(profileAnnotationKeyThrottled),
Value: rawAnnotation,
},
}
// First call should process the annotation
processed1, err := convertAnnotation(timedAnnotation, 0)
require.NoError(t, err)
require.NotNil(t, processed1)
// Second call with the same duplicateTracker should skip
processed2, err := convertAnnotation(timedAnnotation, processed1.duplicateTracker)
require.NoError(t, err)
require.Nil(t, processed2)
})
}
func TestProcessAnnotations(t *testing.T) {
rawAnnotation := `{"body":{"periodType":"day","periodLimitMb":1024,"limitResetTime":1609459200}}`
t.Run("processes multiple annotations", func(t *testing.T) {
annotations := []*TimedAnnotation{
{
Timestamp: 1609455600000,
Annotation: &typesv1.ProfileAnnotation{
Key: string(profileAnnotationKeyThrottled),
Value: rawAnnotation,
},
},
{
Timestamp: 1609459200000,
Annotation: &typesv1.ProfileAnnotation{
Key: string(profileAnnotationKeyThrottled),
Value: rawAnnotation,
},
},
}
result, err := processAnnotations(annotations)
require.NoError(t, err)
require.Equal(t, 1, len(result.times))
require.Equal(t, 1, len(result.timeEnds))
require.Equal(t, 1, len(result.texts))
require.Equal(t, 1, len(result.isRegions))
})
t.Run("handles empty annotations list", func(t *testing.T) {
result, err := processAnnotations([]*TimedAnnotation{})
require.NoError(t, err)
require.Equal(t, 0, len(result.times))
require.Equal(t, 0, len(result.timeEnds))
require.Equal(t, 0, len(result.texts))
require.Equal(t, 0, len(result.isRegions))
})
t.Run("handles nil annotations", func(t *testing.T) {
annotations := []*TimedAnnotation{nil}
result, err := processAnnotations(annotations)
require.NoError(t, err)
require.Equal(t, 0, len(result.times))
})
t.Run("handles invalid annotation data", func(t *testing.T) {
annotations := []*TimedAnnotation{
{
Timestamp: 1000,
Annotation: &typesv1.ProfileAnnotation{
Key: string(profileAnnotationKeyThrottled),
Value: `invalid json`,
},
},
}
result, err := processAnnotations(annotations)
require.Error(t, err)
require.Nil(t, result)
require.Contains(t, err.Error(), "error parsing annotation data")
})
}
func TestCreateAnnotationFrame(t *testing.T) {
rawAnnotation := `{"body":{"periodType":"day","periodLimitMb":1024,"limitResetTime":1609459200}}`
t.Run("creates frame with correct fields", func(t *testing.T) {
annotations := []*TimedAnnotation{
{
Timestamp: 1609455600000,
Annotation: &typesv1.ProfileAnnotation{
Key: string(profileAnnotationKeyThrottled),
Value: rawAnnotation,
},
},
}
frame, err := createAnnotationFrame(annotations)
require.NoError(t, err)
require.NotNil(t, frame)
require.Equal(t, "annotations", frame.Name)
require.Equal(t, data.DataTopicAnnotations, frame.Meta.DataTopic)
require.Equal(t, 5, len(frame.Fields))
require.Equal(t, "time", frame.Fields[0].Name)
require.Equal(t, "timeEnd", frame.Fields[1].Name)
require.Equal(t, "text", frame.Fields[2].Name)
require.Equal(t, "isRegion", frame.Fields[3].Name)
require.Equal(t, "color", frame.Fields[4].Name)
require.Equal(t, 1, frame.Fields[0].Len())
require.Equal(t, time.UnixMilli(1609455600000), frame.Fields[0].At(0))
require.Equal(t, time.UnixMilli(1609459200000), frame.Fields[1].At(0))
require.Contains(t, frame.Fields[2].At(0).(string), "Ingestion limit")
})
t.Run("handles empty annotations list", func(t *testing.T) {
frame, err := createAnnotationFrame([]*TimedAnnotation{})
require.NoError(t, err)
require.NotNil(t, frame)
require.Equal(t, 5, len(frame.Fields))
require.Equal(t, 0, frame.Fields[0].Len())
})
}
@@ -41,6 +41,8 @@ type GrafanaPyroscopeDataQuery struct {
// Specify the query flavor
// TODO make this required and give it a default
QueryType *string `json:"queryType,omitempty"`
// If set to true, the response will contain annotations
Annotations *bool `json:"annotations,omitempty"`
// For mixed data sources the selected datasource is on the query level.
// For non mixed scenarios this is undefined.
// TODO find a better way to do this ^ that's friendly to schema
@@ -46,7 +46,8 @@ type LabelPair struct {
type Point struct {
Value float64
// Milliseconds unix timestamp
Timestamp int64
Timestamp int64
Annotations []*typesv1.ProfileAnnotation
}
type ProfileResponse struct {
@@ -133,8 +134,9 @@ func (c *PyroscopeClient) GetSeries(ctx context.Context, profileTypeID string, l
points := make([]*Point, len(s.Points))
for i, p := range s.Points {
points[i] = &Point{
Value: p.Value,
Timestamp: p.Timestamp,
Value: p.Value,
Timestamp: p.Timestamp,
Annotations: p.Annotations,
}
}
+44 -4
View File
@@ -14,6 +14,7 @@ import (
"github.com/grafana/grafana-plugin-sdk-go/data"
"github.com/grafana/grafana-plugin-sdk-go/live"
"github.com/grafana/grafana/pkg/tsdb/grafana-pyroscope-datasource/kinds/dataquery"
typesv1 "github.com/grafana/pyroscope/api/gen/proto/go/types/v1"
"github.com/xlab/treeprint"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes"
@@ -94,7 +95,15 @@ func (d *PyroscopeDatasource) query(ctx context.Context, pCtx backend.PluginCont
}
// add the frames to the response.
responseMutex.Lock()
response.Frames = append(response.Frames, seriesToDataFrames(seriesResp)...)
withAnnotations := qm.Annotations != nil && *qm.Annotations
frames, err := seriesToDataFrames(seriesResp, withAnnotations)
if err != nil {
span.RecordError(err)
span.SetStatus(codes.Error, err.Error())
logger.Error("Querying SelectSeries()", "err", err, "function", logEntrypoint())
return err
}
response.Frames = append(response.Frames, frames...)
responseMutex.Unlock()
return nil
})
@@ -411,8 +420,22 @@ func walkTree(tree *ProfileTree, fn func(tree *ProfileTree)) {
}
}
func seriesToDataFrames(resp *SeriesResponse) []*data.Frame {
type TimedAnnotation struct {
Timestamp int64 `json:"timestamp"`
Annotation *typesv1.ProfileAnnotation `json:"annotation"`
}
func (ta *TimedAnnotation) getKey() string {
return ta.Annotation.Key
}
func (ta *TimedAnnotation) getValue() string {
return ta.Annotation.Value
}
func seriesToDataFrames(resp *SeriesResponse, withAnnotations bool) ([]*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
@@ -430,15 +453,32 @@ func seriesToDataFrames(resp *SeriesResponse) []*data.Frame {
valueField := data.NewField(resp.Label, labels, []float64{})
valueField.Config = &data.FieldConfig{Unit: resp.Units}
fields = append(fields, valueField)
for _, point := range series.Points {
timeField.Append(time.UnixMilli(point.Timestamp))
valueField.Append(point.Value)
if withAnnotations {
for _, a := range point.Annotations {
annotations = append(annotations, &TimedAnnotation{
Timestamp: point.Timestamp,
Annotation: a,
})
}
}
}
fields = append(fields, valueField)
frame.Fields = fields
frames = append(frames, frame)
}
return frames
if len(annotations) > 0 {
frame, err := createAnnotationFrame(annotations)
if err != nil {
return nil, err
}
frames = append(frames, frame)
}
return frames, nil
}
@@ -7,6 +7,7 @@ import (
"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"
)
@@ -226,6 +227,145 @@ func Test_treeToNestedDataFrame(t *testing.T) {
})
}
func Test_seriesToDataFrameAnnotations(t *testing.T) {
t.Run("annotations field is not added when no annotations are present", func(t *testing.T) {
series := &SeriesResponse{
Series: []*Series{
{
Labels: []*LabelPair{},
Points: []*Point{
{
Timestamp: int64(1000),
Value: 30,
},
{
Timestamp: int64(2000),
Value: 20,
},
{
Timestamp: int64(3000),
Value: 10,
},
},
},
},
Units: "short",
Label: "samples",
}
frames, err := seriesToDataFrames(series, true)
require.NoError(t, err)
require.Equal(t, 1, len(frames))
require.Equal(t, 2, len(frames[0].Fields))
})
t.Run("annotations frame can be skipped", func(t *testing.T) {
rawAnnotation := `{"body":{"periodType":"day","periodLimitMb":1024,"limitResetTime":1609459200}}`
series := &SeriesResponse{
Series: []*Series{
{
Points: []*Point{
{
Timestamp: int64(1609455600000),
Value: 30,
Annotations: []*typesv1.ProfileAnnotation{
{Key: string(profileAnnotationKeyThrottled), Value: rawAnnotation},
},
},
},
},
},
}
frames, err := seriesToDataFrames(series, false)
require.NoError(t, err)
require.Equal(t, 1, len(frames))
})
t.Run("throttling annotations are correctly processed", func(t *testing.T) {
rawAnnotation := `{"body":{"periodType":"day","periodLimitMb":1024,"limitResetTime":1609459200}}`
series := &SeriesResponse{
Series: []*Series{
{
Points: []*Point{
{
Timestamp: int64(1609455600000),
Value: 30,
Annotations: []*typesv1.ProfileAnnotation{
{Key: string(profileAnnotationKeyThrottled), Value: rawAnnotation},
},
},
},
},
},
}
frames, err := seriesToDataFrames(series, true)
require.NoError(t, err)
require.Equal(t, 2, len(frames))
annotationsFrame := frames[1]
require.Equal(t, "annotations", annotationsFrame.Name)
require.Equal(t, data.DataTopicAnnotations, annotationsFrame.Meta.DataTopic)
require.Equal(t, 5, len(annotationsFrame.Fields))
require.Equal(t, "time", annotationsFrame.Fields[0].Name)
require.Equal(t, "timeEnd", annotationsFrame.Fields[1].Name)
require.Equal(t, "text", annotationsFrame.Fields[2].Name)
require.Equal(t, "isRegion", annotationsFrame.Fields[3].Name)
require.Equal(t, "color", annotationsFrame.Fields[4].Name)
require.Equal(t, 1, annotationsFrame.Fields[0].Len())
require.Equal(t, time.UnixMilli(1609455600000), annotationsFrame.Fields[0].At(0))
require.Equal(t, time.UnixMilli(1609459200000), annotationsFrame.Fields[1].At(0))
require.Contains(t, annotationsFrame.Fields[2].At(0).(string), "Ingestion limit")
})
t.Run("non-throttling annotations are ignored", func(t *testing.T) {
series := &SeriesResponse{
Series: []*Series{
{
Points: []*Point{
{
Timestamp: int64(1000),
Value: 30,
Annotations: []*typesv1.ProfileAnnotation{
{Key: "key1", Value: "value1"},
{Key: "key2", Value: "value2"},
},
},
{
Timestamp: int64(2000),
Value: 20,
Annotations: []*typesv1.ProfileAnnotation{
{Key: "key3", Value: "value3"},
},
},
},
},
},
}
frames, err := seriesToDataFrames(series, true)
require.NoError(t, err)
require.Equal(t, 2, len(frames))
annotationsFrame := frames[1]
require.Equal(t, "annotations", annotationsFrame.Name)
require.Equal(t, data.DataTopicAnnotations, annotationsFrame.Meta.DataTopic)
require.Equal(t, 5, len(annotationsFrame.Fields))
require.Equal(t, 0, annotationsFrame.Fields[0].Len())
require.Equal(t, 0, annotationsFrame.Fields[1].Len())
require.Equal(t, 0, annotationsFrame.Fields[2].Len())
require.Equal(t, 0, annotationsFrame.Fields[3].Len())
require.Equal(t, 0, annotationsFrame.Fields[4].Len())
})
}
func Test_seriesToDataFrame(t *testing.T) {
t.Run("single series", func(t *testing.T) {
series := &SeriesResponse{
@@ -235,7 +375,8 @@ func Test_seriesToDataFrame(t *testing.T) {
Units: "short",
Label: "samples",
}
frames := seriesToDataFrames(series)
frames, err := seriesToDataFrames(series, true)
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])
require.Equal(t, data.NewField("samples", map[string]string{}, []float64{30, 10}).SetConfig(&data.FieldConfig{Unit: "short"}), frames[0].Fields[1])
@@ -249,7 +390,8 @@ func Test_seriesToDataFrame(t *testing.T) {
Label: "samples",
}
frames = seriesToDataFrames(series)
frames, err = seriesToDataFrames(series, true)
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])
})
@@ -262,7 +404,8 @@ func Test_seriesToDataFrame(t *testing.T) {
Units: "short",
Label: "samples",
}
frames := seriesToDataFrames(resp)
frames, err := seriesToDataFrames(resp, true)
require.NoError(t, err)
require.Equal(t, 2, len(frames))
require.Equal(t, 2, len(frames[0].Fields))
require.Equal(t, 2, len(frames[1].Fields))