-
Notifications
You must be signed in to change notification settings - Fork 578
/
frontend_select_merge_span_profile.go
89 lines (78 loc) · 3.1 KB
/
frontend_select_merge_span_profile.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
package frontend
import (
"context"
"time"
"connectrpc.com/connect"
"github.com/grafana/dskit/tenant"
"github.com/opentracing/opentracing-go"
"github.com/prometheus/common/model"
"golang.org/x/sync/errgroup"
querierv1 "github.com/grafana/pyroscope/api/gen/proto/go/querier/v1"
"github.com/grafana/pyroscope/api/gen/proto/go/querier/v1/querierv1connect"
phlaremodel "github.com/grafana/pyroscope/pkg/model"
"github.com/grafana/pyroscope/pkg/util/connectgrpc"
validationutil "github.com/grafana/pyroscope/pkg/util/validation"
"github.com/grafana/pyroscope/pkg/validation"
)
func (f *Frontend) SelectMergeSpanProfile(ctx context.Context,
c *connect.Request[querierv1.SelectMergeSpanProfileRequest]) (
*connect.Response[querierv1.SelectMergeSpanProfileResponse], error,
) {
opentracing.SpanFromContext(ctx).
SetTag("start", model.Time(c.Msg.Start).Time().String()).
SetTag("end", model.Time(c.Msg.End).Time().String()).
SetTag("selector", c.Msg.LabelSelector).
SetTag("max_nodes", c.Msg.MaxNodes).
SetTag("profile_type", c.Msg.ProfileTypeID)
ctx = connectgrpc.WithProcedure(ctx, querierv1connect.QuerierServiceSelectMergeSpanProfileProcedure)
tenantIDs, err := tenant.TenantIDs(ctx)
if err != nil {
return nil, connect.NewError(connect.CodeInvalidArgument, err)
}
validated, err := validation.ValidateRangeRequest(f.limits, tenantIDs, model.Interval{Start: model.Time(c.Msg.Start), End: model.Time(c.Msg.End)}, model.Now())
if err != nil {
return nil, connect.NewError(connect.CodeInvalidArgument, err)
}
if validated.IsEmpty {
return connect.NewResponse(&querierv1.SelectMergeSpanProfileResponse{Flamegraph: &querierv1.FlameGraph{}}), nil
}
maxNodes, err := validation.ValidateMaxNodes(f.limits, tenantIDs, c.Msg.GetMaxNodes())
if err != nil {
return nil, connect.NewError(connect.CodeInvalidArgument, err)
}
g, ctx := errgroup.WithContext(ctx)
if maxConcurrent := validationutil.SmallestPositiveNonZeroIntPerTenant(tenantIDs, f.limits.MaxQueryParallelism); maxConcurrent > 0 {
g.SetLimit(maxConcurrent)
}
m := phlaremodel.NewFlameGraphMerger()
interval := validationutil.MaxDurationOrZeroPerTenant(tenantIDs, f.limits.QuerySplitDuration)
intervals := NewTimeIntervalIterator(time.UnixMilli(int64(validated.Start)), time.UnixMilli(int64(validated.End)), interval)
for intervals.Next() {
r := intervals.At()
g.Go(func() error {
req := connectgrpc.CloneRequest(c, &querierv1.SelectMergeSpanProfileRequest{
ProfileTypeID: c.Msg.ProfileTypeID,
LabelSelector: c.Msg.LabelSelector,
Start: r.Start.UnixMilli(),
End: r.End.UnixMilli(),
MaxNodes: &maxNodes,
SpanSelector: c.Msg.SpanSelector,
})
resp, err := connectgrpc.RoundTripUnary[
querierv1.SelectMergeSpanProfileRequest,
querierv1.SelectMergeSpanProfileResponse](ctx, f, req)
if err != nil {
return err
}
m.MergeFlameGraph(resp.Msg.Flamegraph)
return nil
})
}
if err = g.Wait(); err != nil {
return nil, err
}
t := m.Tree()
return connect.NewResponse(&querierv1.SelectMergeSpanProfileResponse{
Flamegraph: phlaremodel.NewFlameGraph(t, c.Msg.GetMaxNodes()),
}), nil
}