From e3f5a6537258dc5e8e5f613dcbfb050275881392 Mon Sep 17 00:00:00 2001 From: Aleksandar Petrov <8142643+aleks-p@users.noreply.github.com> Date: Fri, 29 Aug 2025 08:14:22 -0300 Subject: [PATCH] Pyroscope: Process and display sampling annotations (#109707) * pyroscope: process sampling annotations * Enable annotations in classic explore * Run prettier * Revert unneeded change to plugin.json * Tweak wording in sampling annotation * Fix test * Disable annotations by default --- .../annotation/annotation.go | 123 ++++++ .../annotation/annotation_test.go | 355 ++++++++++++++++++ .../annotation/sampling.go | 47 +++ .../annotation/throttling.go | 42 +++ .../annotations.go | 133 ------- .../annotations_test.go | 188 ---------- .../grafana-pyroscope-datasource/query.go | 24 +- .../query_test.go | 6 +- .../QueryEditor/QueryOptions.tsx | 10 +- 9 files changed, 586 insertions(+), 342 deletions(-) create mode 100644 pkg/tsdb/grafana-pyroscope-datasource/annotation/annotation.go create mode 100644 pkg/tsdb/grafana-pyroscope-datasource/annotation/annotation_test.go create mode 100644 pkg/tsdb/grafana-pyroscope-datasource/annotation/sampling.go create mode 100644 pkg/tsdb/grafana-pyroscope-datasource/annotation/throttling.go delete mode 100644 pkg/tsdb/grafana-pyroscope-datasource/annotations.go delete mode 100644 pkg/tsdb/grafana-pyroscope-datasource/annotations_test.go diff --git a/pkg/tsdb/grafana-pyroscope-datasource/annotation/annotation.go b/pkg/tsdb/grafana-pyroscope-datasource/annotation/annotation.go new file mode 100644 index 00000000000..186007b5bc4 --- /dev/null +++ b/pkg/tsdb/grafana-pyroscope-datasource/annotation/annotation.go @@ -0,0 +1,123 @@ +package annotation + +import ( + "time" + + "github.com/grafana/grafana-plugin-sdk-go/data" + typesv1 "github.com/grafana/pyroscope/api/gen/proto/go/types/v1" +) + +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 +} + +type profileAnnotationKey string + +const ( + // ProfileAnnotationKeyThrottled is an identifier for throttling annotations + ProfileAnnotationKeyThrottled profileAnnotationKey = "pyroscope.ingest.throttled" + // ProfileAnnotationKeySampled is an identifier for sampling annotations + ProfileAnnotationKeySampled profileAnnotationKey = "pyroscope.ingest.sampled" +) + +type processedProfileAnnotation struct { + id string + text string + time int64 + timeEnd int64 + isRegion bool +} + +type grafanaAnnotationData struct { + ids []string + times []time.Time + timeEnds []time.Time + texts []string + isRegions []bool +} + +func (ga *grafanaAnnotationData) add(a *processedProfileAnnotation) { + // simple de-duplication, assuming annotations are ordered + if len(ga.ids) > 0 { + lastIdx := len(ga.ids) - 1 + if a.id == ga.ids[lastIdx] { + // duplicate annotation, extend the previous annotation and discard the rest + ga.timeEnds[lastIdx] = time.UnixMilli(a.timeEnd) + return + } + } + ga.ids = append(ga.ids, a.id) + ga.times = append(ga.times, time.UnixMilli(a.time)) + ga.timeEnds = append(ga.timeEnds, time.UnixMilli(a.timeEnd)) + ga.isRegions = append(ga.isRegions, a.isRegion) + ga.texts = append(ga.texts, a.text) +} + +// convertAnnotation converts a Pyroscope profile annotation into a Grafana annotation +func convertAnnotation(timedAnnotation *TimedAnnotation) (*processedProfileAnnotation, error) { + switch timedAnnotation.getKey() { + case string(ProfileAnnotationKeySampled): + return convertSamplingAnnotation(timedAnnotation.getValue(), timedAnnotation.Timestamp) + case string(ProfileAnnotationKeyThrottled): + return convertThrottlingAnnotation(timedAnnotation.getValue(), timedAnnotation.Timestamp) + default: + // Currently, we only support throttling and sampling annotations + return nil, nil + } +} + +func processAnnotations(timedAnnotations []*TimedAnnotation) (*grafanaAnnotationData, error) { + result := &grafanaAnnotationData{ + times: []time.Time{}, + timeEnds: []time.Time{}, + texts: []string{}, + isRegions: []bool{}, + } + + for _, timedAnnotation := range timedAnnotations { + if timedAnnotation == nil || timedAnnotation.Annotation == nil { + continue + } + processed, err := convertAnnotation(timedAnnotation) + if err != nil { + return nil, err + } + + if processed != nil { + result.add(processed) + } + } + + return result, nil +} + +// CreateAnnotationFrame creates a Grafana data frame from annotation data +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 +} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/annotation/annotation_test.go b/pkg/tsdb/grafana-pyroscope-datasource/annotation/annotation_test.go new file mode 100644 index 00000000000..0d9f9e72d29 --- /dev/null +++ b/pkg/tsdb/grafana-pyroscope-datasource/annotation/annotation_test.go @@ -0,0 +1,355 @@ +package annotation + +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) { + t.Run("converts a valid throttling annotation", func(t *testing.T) { + rawAnnotation := `{"body":{"periodType":"day","periodLimitMb":1024,"limitResetTime":1609459200}}` + timedAnnotation := &TimedAnnotation{ + Timestamp: 1609455600000, + Annotation: &typesv1.ProfileAnnotation{ + Key: string(ProfileAnnotationKeyThrottled), + Value: rawAnnotation, + }, + } + + processed, err := convertAnnotation(timedAnnotation) + 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 + }) + + t.Run("converts a valid sampling annotation", func(t *testing.T) { + rawAnnotation := `{"body":{"source": {"usageGroup":"group-1","probability":0.1}}}` + timedAnnotation := &TimedAnnotation{ + Timestamp: 1609455600000, + Annotation: &typesv1.ProfileAnnotation{ + Key: string(ProfileAnnotationKeySampled), + Value: rawAnnotation, + }, + } + + processed, err := convertAnnotation(timedAnnotation) + require.NoError(t, err) + require.NotNil(t, processed) + require.Contains(t, processed.text, "Profile volume reduced by 90.00% for this service.") + require.Equal(t, int64(1609455600000), processed.time) + require.Equal(t, int64(1609455600000), processed.timeEnd) + }) + + 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) + 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) + require.Error(t, err) + require.Nil(t, processed) + require.Contains(t, err.Error(), "error parsing annotation data") + }) +} + +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 TestGrafanaAnnotationDataAdd(t *testing.T) { + t.Run("adds first annotation", func(t *testing.T) { + ga := &grafanaAnnotationData{ + ids: []string{}, + times: []time.Time{}, + timeEnds: []time.Time{}, + texts: []string{}, + isRegions: []bool{}, + } + + annotation := &processedProfileAnnotation{ + id: "test-id-1", + text: "Test annotation 1", + time: 1609455600000, + timeEnd: 1609459200000, + isRegion: true, + } + + ga.add(annotation) + + require.Equal(t, 1, len(ga.ids)) + require.Equal(t, "test-id-1", ga.ids[0]) + require.Equal(t, time.UnixMilli(1609455600000), ga.times[0]) + require.Equal(t, time.UnixMilli(1609459200000), ga.timeEnds[0]) + require.Equal(t, "Test annotation 1", ga.texts[0]) + require.Equal(t, true, ga.isRegions[0]) + }) + + t.Run("adds different annotations", func(t *testing.T) { + ga := &grafanaAnnotationData{ + ids: []string{}, + times: []time.Time{}, + timeEnds: []time.Time{}, + texts: []string{}, + isRegions: []bool{}, + } + + annotation1 := &processedProfileAnnotation{ + id: "test-id-1", + text: "Test annotation 1", + time: 1609455600000, + timeEnd: 1609459200000, + isRegion: true, + } + + annotation2 := &processedProfileAnnotation{ + id: "test-id-2", + text: "Test annotation 2", + time: 1609463800000, + timeEnd: 1609467400000, + isRegion: false, + } + + ga.add(annotation1) + ga.add(annotation2) + + require.Equal(t, 2, len(ga.ids)) + require.Equal(t, "test-id-1", ga.ids[0]) + require.Equal(t, "test-id-2", ga.ids[1]) + require.Equal(t, time.UnixMilli(1609455600000), ga.times[0]) + require.Equal(t, time.UnixMilli(1609463800000), ga.times[1]) + }) + + t.Run("removes duplicates and extends timeEnd", func(t *testing.T) { + ga := &grafanaAnnotationData{ + ids: []string{}, + times: []time.Time{}, + timeEnds: []time.Time{}, + texts: []string{}, + isRegions: []bool{}, + } + + annotation1 := &processedProfileAnnotation{ + id: "duplicate-id", + text: "First occurrence", + time: 1609455600000, + timeEnd: 1609459200000, + isRegion: true, + } + + annotation2 := &processedProfileAnnotation{ + id: "duplicate-id", + text: "Second occurrence (should be ignored)", + time: 1609460000000, + timeEnd: 1609463600000, + isRegion: false, + } + + ga.add(annotation1) + ga.add(annotation2) + + require.Equal(t, 1, len(ga.ids)) + require.Equal(t, 1, len(ga.times)) + require.Equal(t, 1, len(ga.timeEnds)) + require.Equal(t, 1, len(ga.texts)) + require.Equal(t, 1, len(ga.isRegions)) + + require.Equal(t, "duplicate-id", ga.ids[0]) + require.Equal(t, time.UnixMilli(1609455600000), ga.times[0]) // Original time + require.Equal(t, time.UnixMilli(1609463600000), ga.timeEnds[0]) // Extended timeEnd + require.Equal(t, "First occurrence", ga.texts[0]) // Original text + require.Equal(t, true, ga.isRegions[0]) // Original isRegion + }) + + t.Run("handles multiple duplicates correctly", func(t *testing.T) { + ga := &grafanaAnnotationData{ + ids: []string{}, + times: []time.Time{}, + timeEnds: []time.Time{}, + texts: []string{}, + isRegions: []bool{}, + } + + annotation1 := &processedProfileAnnotation{ + id: "id-1", + text: "Annotation 1", + time: 1609455600000, + timeEnd: 1609459200000, + isRegion: true, + } + + // Add duplicate of first + annotation1Duplicate := &processedProfileAnnotation{ + id: "id-1", + text: "Annotation 1 duplicate", + time: 1609460000000, + timeEnd: 1609470000000, + isRegion: false, + } + + // Add a second, unique annotation + annotation2 := &processedProfileAnnotation{ + id: "id-2", + text: "Annotation 2", + time: 1609480000000, + timeEnd: 1609490000000, + isRegion: false, + } + + // Add duplicate of second + annotation2Duplicate := &processedProfileAnnotation{ + id: "id-2", + text: "Annotation 2 duplicate", + time: 1609500000000, + timeEnd: 1609510000000, + isRegion: true, + } + + ga.add(annotation1) + ga.add(annotation1Duplicate) + ga.add(annotation2) + ga.add(annotation2Duplicate) + + require.Equal(t, 2, len(ga.ids)) + require.Equal(t, "id-1", ga.ids[0]) + require.Equal(t, "id-2", ga.ids[1]) + + // The first annotation should have an extended timeEnd + require.Equal(t, time.UnixMilli(1609455600000), ga.times[0]) + require.Equal(t, time.UnixMilli(1609470000000), ga.timeEnds[0]) + require.Equal(t, "Annotation 1", ga.texts[0]) + require.Equal(t, true, ga.isRegions[0]) + + // The second annotation should have an extended timeEnd + require.Equal(t, time.UnixMilli(1609480000000), ga.times[1]) + require.Equal(t, time.UnixMilli(1609510000000), ga.timeEnds[1]) + require.Equal(t, "Annotation 2", ga.texts[1]) + require.Equal(t, false, ga.isRegions[1]) + }) +} + +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()) + }) +} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/annotation/sampling.go b/pkg/tsdb/grafana-pyroscope-datasource/annotation/sampling.go new file mode 100644 index 00000000000..49121163f61 --- /dev/null +++ b/pkg/tsdb/grafana-pyroscope-datasource/annotation/sampling.go @@ -0,0 +1,47 @@ +package annotation + +import ( + "encoding/json" + "fmt" +) + +type annotationWithSamplingBody struct { + Body profileSampledAnnotation `json:"body"` +} + +type profileSampledAnnotation struct { + Source *samplingSource `json:"source"` +} + +type samplingSource struct { + UsageGroup string `json:"usageGroup"` + Probability float64 `json:"probability"` +} + +func convertSamplingAnnotation(raw string, timestamp int64) (*processedProfileAnnotation, error) { + var profileAnnotation annotationWithSamplingBody + err := json.Unmarshal([]byte(raw), &profileAnnotation) + if err != nil { + return nil, fmt.Errorf("error parsing annotation data: %w", err) + } + if profileAnnotation.Body.Source == nil { + return nil, fmt.Errorf("error parsing sampling annotation data: source is nil") + } + + samplingInfo := profileAnnotation.Body.Source + if samplingInfo.Probability == 1.0 { + return nil, nil + } + + reductionPercentage := (1 - samplingInfo.Probability) * 100 + id := fmt.Sprintf("%s-%.0f", samplingInfo.UsageGroup, reductionPercentage) + text := fmt.Sprintf("Profile volume reduced by %.2f%% for this service.", reductionPercentage) + + return &processedProfileAnnotation{ + id: id, + text: text, + time: timestamp, + timeEnd: timestamp, + isRegion: true, + }, nil +} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/annotation/throttling.go b/pkg/tsdb/grafana-pyroscope-datasource/annotation/throttling.go new file mode 100644 index 00000000000..151cb30b30e --- /dev/null +++ b/pkg/tsdb/grafana-pyroscope-datasource/annotation/throttling.go @@ -0,0 +1,42 @@ +package annotation + +import ( + "encoding/json" + "fmt" + "time" + + "github.com/dustin/go-humanize" +) + +type annotationWithThrottlingBody struct { + Body profileThrottledAnnotation `json:"body"` +} + +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"` +} + +func convertThrottlingAnnotation(raw string, timestamp int64) (*processedProfileAnnotation, error) { + var profileAnnotation annotationWithThrottlingBody + err := json.Unmarshal([]byte(raw), &profileAnnotation) + if err != nil { + return nil, fmt.Errorf("error parsing annotation data: %w", err) + } + + throttlingInfo := profileAnnotation.Body + limit := humanize.IBytes(uint64(throttlingInfo.PeriodLimitMb * 1024 * 1024)) + id := fmt.Sprintf("%s-%s-%d", throttlingInfo.PeriodType, limit, throttlingInfo.LimitResetTime) + + return &processedProfileAnnotation{ + id: id, + text: fmt.Sprintf("Ingestion limit (%s/%s) reached", limit, throttlingInfo.PeriodType), + time: timestamp, + timeEnd: throttlingInfo.LimitResetTime * 1000, + isRegion: throttlingInfo.LimitResetTime < time.Now().Unix(), + }, nil +} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/annotations.go b/pkg/tsdb/grafana-pyroscope-datasource/annotations.go deleted file mode 100644 index aac2778895d..00000000000 --- a/pkg/tsdb/grafana-pyroscope-datasource/annotations.go +++ /dev/null @@ -1,133 +0,0 @@ -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 -} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/annotations_test.go b/pkg/tsdb/grafana-pyroscope-datasource/annotations_test.go deleted file mode 100644 index 3ca7d7433f1..00000000000 --- a/pkg/tsdb/grafana-pyroscope-datasource/annotations_test.go +++ /dev/null @@ -1,188 +0,0 @@ -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()) - }) -} diff --git a/pkg/tsdb/grafana-pyroscope-datasource/query.go b/pkg/tsdb/grafana-pyroscope-datasource/query.go index 506f9de33ee..88d1d216956 100644 --- a/pkg/tsdb/grafana-pyroscope-datasource/query.go +++ b/pkg/tsdb/grafana-pyroscope-datasource/query.go @@ -13,13 +13,14 @@ import ( "github.com/grafana/grafana-plugin-sdk-go/backend/tracing" "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" "go.opentelemetry.io/otel/trace" "golang.org/x/sync/errgroup" + + "github.com/grafana/grafana/pkg/tsdb/grafana-pyroscope-datasource/annotation" + "github.com/grafana/grafana/pkg/tsdb/grafana-pyroscope-datasource/kinds/dataquery" ) type queryModel struct { @@ -454,19 +455,6 @@ func walkTree(tree *ProfileTree, fn func(tree *ProfileTree)) { } } -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 -} - // isCumulativeProfile determines if a profile type requires rate calculation using the metadata registry func isCumulativeProfile(profileTypeID string) bool { registry := GetProfileMetadataRegistry() @@ -500,7 +488,7 @@ func convertToRateUnit(originalUnit string) string { 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) + annotations := make([]*annotation.TimedAnnotation, 0) for _, series := range resp.Series { // We create separate data frames as the series may not have the same length @@ -555,7 +543,7 @@ func seriesToDataFrames(resp *SeriesResponse, withAnnotations bool, stepDuration valueField.Append(value) if withAnnotations { for _, a := range point.Annotations { - annotations = append(annotations, &TimedAnnotation{ + annotations = append(annotations, &annotation.TimedAnnotation{ Timestamp: point.Timestamp, Annotation: a, }) @@ -568,7 +556,7 @@ func seriesToDataFrames(resp *SeriesResponse, withAnnotations bool, stepDuration } if len(annotations) > 0 { - frame, err := createAnnotationFrame(annotations) + frame, err := annotation.CreateAnnotationFrame(annotations) if err != nil { return nil, err } diff --git a/pkg/tsdb/grafana-pyroscope-datasource/query_test.go b/pkg/tsdb/grafana-pyroscope-datasource/query_test.go index 67addaaedd4..1f1404da99f 100644 --- a/pkg/tsdb/grafana-pyroscope-datasource/query_test.go +++ b/pkg/tsdb/grafana-pyroscope-datasource/query_test.go @@ -10,6 +10,8 @@ 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/grafana/grafana/pkg/tsdb/grafana-pyroscope-datasource/annotation" ) // This is where the tests for the datasource backend live. @@ -313,7 +315,7 @@ func Test_seriesToDataFrameAnnotations(t *testing.T) { Timestamp: int64(1609455600000), Value: 30, Annotations: []*typesv1.ProfileAnnotation{ - {Key: string(profileAnnotationKeyThrottled), Value: rawAnnotation}, + {Key: string(annotation.ProfileAnnotationKeyThrottled), Value: rawAnnotation}, }, }, }, @@ -337,7 +339,7 @@ func Test_seriesToDataFrameAnnotations(t *testing.T) { Timestamp: int64(1609455600000), Value: 30, Annotations: []*typesv1.ProfileAnnotation{ - {Key: string(profileAnnotationKeyThrottled), Value: rawAnnotation}, + {Key: string(annotation.ProfileAnnotationKeyThrottled), Value: rawAnnotation}, }, }, }, diff --git a/public/app/plugins/datasource/grafana-pyroscope-datasource/QueryEditor/QueryOptions.tsx b/public/app/plugins/datasource/grafana-pyroscope-datasource/QueryEditor/QueryOptions.tsx index 5b2774f3105..1ae0dca0241 100644 --- a/public/app/plugins/datasource/grafana-pyroscope-datasource/QueryEditor/QueryOptions.tsx +++ b/public/app/plugins/datasource/grafana-pyroscope-datasource/QueryEditor/QueryOptions.tsx @@ -2,7 +2,7 @@ import { css } from '@emotion/css'; import * as React from 'react'; import { CoreApp, GrafanaTheme2, SelectableValue } from '@grafana/data'; -import { useStyles2, RadioButtonGroup, MultiSelect, Input } from '@grafana/ui'; +import { useStyles2, RadioButtonGroup, MultiSelect, Input, InlineSwitch } from '@grafana/ui'; import { Query } from '../types'; @@ -134,6 +134,14 @@ export function QueryOptions({ query, onQueryChange, app, labels }: Props) { }} /> + Include profiling annotations in the time series.}> + ) => { + onQueryChange({ ...query, annotations: event.currentTarget.checked }); + }} + /> +