Skip to content

Commit dd039bd

Browse files
author
morvencao
committed
split grpc common metrics and cloudevents metrics.
Signed-off-by: morvencao <lcao@redhat.com> rh-pre-commit.version: 2.3.2 rh-pre-commit.check-secrets: ENABLED
1 parent cc0f990 commit dd039bd

6 files changed

Lines changed: 322 additions & 276 deletions

File tree

Lines changed: 306 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,306 @@
1+
package metrics
2+
3+
import (
4+
"context"
5+
"fmt"
6+
"sync"
7+
"time"
8+
9+
"google.golang.org/grpc"
10+
"google.golang.org/grpc/status"
11+
12+
"k8s.io/component-base/metrics/legacyregistry"
13+
"k8s.io/klog/v2"
14+
15+
"github.com/cloudevents/sdk-go/v2/binding"
16+
cetypes "github.com/cloudevents/sdk-go/v2/types"
17+
pbv1 "open-cluster-management.io/sdk-go/pkg/cloudevents/generic/options/grpc/protobuf/v1"
18+
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/options/grpc/protocol"
19+
"open-cluster-management.io/sdk-go/pkg/cloudevents/generic/types"
20+
21+
k8smetrics "k8s.io/component-base/metrics"
22+
)
23+
24+
// ensure metrics are registered only once
25+
var once sync.Once
26+
27+
// subsystem used to define the metrics of grpc server for cloud events
28+
const grpcCEMetricsSubsystem = "grpc_server_ce"
29+
30+
// Names of the labels added to metrics:
31+
const (
32+
grpcCEMetricsClusterLabel = "cluster"
33+
grpcCEMetricsDataTypeLabel = "data_type"
34+
grpcCEMetricsMethodLabel = "method"
35+
grpcMetricsCodeLabel = "grpc_code"
36+
)
37+
38+
// grpcCEMetricsLabels - Array of labels added to grpc server metrics for cloudevents:
39+
var grpcCEMetricsLabels = []string{
40+
grpcCEMetricsClusterLabel,
41+
grpcCEMetricsDataTypeLabel,
42+
grpcCEMetricsMethodLabel,
43+
}
44+
45+
// grpcCEMetricsAllLabels - Array of all labels added to grpc server metrics for cloudevents:
46+
var grpcCEMetricsAllLabels = append(grpcCEMetricsLabels, grpcMetricsCodeLabel)
47+
48+
// Names of the grpc server metrics for cloudevents:
49+
const (
50+
calledCountMetric = "called_total"
51+
processedCountMetric = "processed_total"
52+
processingDurationMetric = "processing_duration_seconds"
53+
messageReceivedCountMetric = "msg_received_total"
54+
messageSentCountMetric = "msg_sent_total"
55+
)
56+
57+
// grpcCECalledCountMetric is a counter metric that tracks the total number of
58+
// RPC requests for cloudevents called on the gRPC server.
59+
var grpcCECalledCountMetric = k8smetrics.NewCounterVec(&k8smetrics.CounterOpts{
60+
Subsystem: grpcCEMetricsSubsystem,
61+
Name: calledCountMetric,
62+
StabilityLevel: k8smetrics.ALPHA,
63+
Help: "Total number of RPC requests for cloudevents called on the grpc server.",
64+
}, grpcCEMetricsLabels)
65+
66+
// grpcCEMessageReceivedCountMetric is a counter metric that tracks the total number of
67+
// messages for cloudevents received on the gRPC server.
68+
var grpcCEMessageReceivedCountMetric = k8smetrics.NewCounterVec(&k8smetrics.CounterOpts{
69+
Subsystem: grpcCEMetricsSubsystem,
70+
Name: messageReceivedCountMetric,
71+
StabilityLevel: k8smetrics.ALPHA,
72+
Help: "Total number of messages for cloudevents received on the gRPC server.",
73+
}, grpcCEMetricsLabels)
74+
75+
// grpcCEMessageSentCountMetric is a counter metric that tracks the total number of
76+
// messages for cloudevents sent by the gRPC server.
77+
var grpcCEMessageSentCountMetric = k8smetrics.NewCounterVec(&k8smetrics.CounterOpts{
78+
Subsystem: grpcCEMetricsSubsystem,
79+
Name: messageSentCountMetric,
80+
StabilityLevel: k8smetrics.ALPHA,
81+
Help: "Total number of messages for cloudevents sent by the gRPC server.",
82+
}, grpcCEMetricsLabels)
83+
84+
// grpcCEProcessedCountMetric is a counter metric that tracks the total number of
85+
// RPC requests for cloudevents processed on the server, regardless of success or failure.
86+
var grpcCEProcessedCountMetric = k8smetrics.NewCounterVec(&k8smetrics.CounterOpts{
87+
Subsystem: grpcCEMetricsSubsystem,
88+
Name: processedCountMetric,
89+
StabilityLevel: k8smetrics.ALPHA,
90+
Help: "Total number of RPC requests for cloudevents processed on the server, regardless of success or failure.",
91+
}, grpcCEMetricsAllLabels,
92+
)
93+
94+
// grpcCEProcessingDurationMetric is a histogram metric that tracks the duration of
95+
// RPC requests for cloudevents processed on the server.
96+
var grpcCEProcessingDurationMetric = k8smetrics.NewHistogramVec(&k8smetrics.HistogramOpts{
97+
Subsystem: grpcCEMetricsSubsystem,
98+
Name: processingDurationMetric,
99+
StabilityLevel: k8smetrics.ALPHA,
100+
Help: "Histogram of the duration of RPC requests for cloudevents processed on the server.",
101+
Buckets: k8smetrics.ExponentialBuckets(10e-7, 10, 10),
102+
}, grpcCEMetricsAllLabels)
103+
104+
const gRPCCloudEventService = "io.cloudevents.v1.CloudEventService"
105+
106+
// NewCloudEventsMetricsUnaryInterceptor creates a unary server interceptor for cloudevents metrics.
107+
func NewCloudEventsMetricsUnaryInterceptor() grpc.UnaryServerInterceptor {
108+
return func(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
109+
startTime := time.Now()
110+
111+
// extract the service and method
112+
service, method := SplitMethod(info.FullMethod)
113+
if service != gRPCCloudEventService {
114+
// skip emitting cloudevents metrics if it's not a cloudevents RPC request
115+
return handler(ctx, req)
116+
}
117+
118+
// initialize defaults for error cases
119+
cluster := "unknown"
120+
dataType := "unknown"
121+
122+
pubReq, ok := req.(*pbv1.PublishRequest)
123+
if !ok {
124+
err := fmt.Errorf("invalid request type for Publish method")
125+
recordCloudEventsMetrics(cluster, dataType, method, err, startTime)
126+
return nil, err
127+
}
128+
// convert the request to cloudevent and extract the source
129+
evt, err := binding.ToEvent(ctx, protocol.NewMessage(pubReq.Event))
130+
if err != nil {
131+
err = fmt.Errorf("failed to convert to cloudevent: %v", err)
132+
recordCloudEventsMetrics(cluster, dataType, method, err, startTime)
133+
return nil, err
134+
}
135+
136+
// extract the cluster name from event extensions
137+
clusterVal, err := cetypes.ToString(evt.Context.GetExtensions()[types.ExtensionClusterName])
138+
if err != nil {
139+
err = fmt.Errorf("failed to get clustername extension: %v", err)
140+
recordCloudEventsMetrics(cluster, dataType, method, err, startTime)
141+
return nil, err
142+
}
143+
cluster = clusterVal
144+
145+
// extract the data type from event type
146+
eventType, err := types.ParseCloudEventsType(evt.Type())
147+
if err != nil {
148+
err = fmt.Errorf("failed to parse cloud event type %s, %v", evt.Type(), err)
149+
recordCloudEventsMetrics(cluster, dataType, method, err, startTime)
150+
return nil, err
151+
}
152+
dataType = eventType.CloudEventsDataType.String()
153+
154+
grpcCECalledCountMetric.WithLabelValues(cluster, dataType, method).Inc()
155+
grpcCEMessageReceivedCountMetric.WithLabelValues(cluster, dataType, method).Inc()
156+
// call rpc handler to handle RPC request
157+
resp, err := handler(ctx, req)
158+
duration := time.Since(startTime).Seconds()
159+
grpcCEMessageSentCountMetric.WithLabelValues(cluster, dataType, method).Inc()
160+
161+
// get status code from error
162+
status := statusFromError(err)
163+
code := status.Code()
164+
grpcCEProcessedCountMetric.WithLabelValues(cluster, dataType, method, code.String()).Inc()
165+
grpcCEProcessingDurationMetric.WithLabelValues(cluster, dataType, method, code.String()).Observe(duration)
166+
167+
return resp, err
168+
}
169+
}
170+
171+
func recordCloudEventsMetrics(cluster, dataType, method string, err error, startTime time.Time) {
172+
duration := time.Since(startTime).Seconds()
173+
status := statusFromError(err)
174+
code := status.Code()
175+
grpcCEProcessedCountMetric.WithLabelValues(cluster, dataType, method, code.String()).Inc()
176+
grpcCEProcessingDurationMetric.WithLabelValues(cluster, dataType, method, code.String()).Observe(duration)
177+
}
178+
179+
// wrappedCloudEventsMetricsStream wraps a grpc.ServerStream, capturing the request source
180+
// emitting metrics for the stream interceptor.
181+
type wrappedCloudEventsMetricsStream struct {
182+
clusterName *string
183+
dataType *string
184+
method string
185+
grpc.ServerStream
186+
ctx context.Context
187+
}
188+
189+
// RecvMsg wraps the RecvMsg method of the embedded grpc.ServerStream.
190+
// It captures the cluster and data type from the SubscriptionRequest and emits metrics.
191+
func (w *wrappedCloudEventsMetricsStream) RecvMsg(m interface{}) error {
192+
err := w.ServerStream.RecvMsg(m)
193+
if err != nil {
194+
return err
195+
}
196+
197+
subReq, ok := m.(*pbv1.SubscriptionRequest)
198+
if !ok {
199+
return fmt.Errorf("invalid request type for Subscribe method")
200+
}
201+
202+
if w.clusterName != nil && w.dataType != nil {
203+
*w.clusterName = subReq.ClusterName
204+
*w.dataType = subReq.DataType
205+
grpcCECalledCountMetric.WithLabelValues(*w.clusterName, *w.dataType, w.method).Inc()
206+
grpcCEMessageReceivedCountMetric.WithLabelValues(*w.clusterName, *w.dataType, w.method).Inc()
207+
}
208+
209+
return nil
210+
}
211+
212+
// SendMsg wraps the SendMsg method of the embedded grpc.ServerStream.
213+
func (w *wrappedCloudEventsMetricsStream) SendMsg(m interface{}) error {
214+
err := w.ServerStream.SendMsg(m)
215+
if err != nil {
216+
return err
217+
}
218+
219+
if w.clusterName != nil && w.dataType != nil && *w.clusterName != "" && *w.dataType != "" {
220+
grpcCEMessageSentCountMetric.WithLabelValues(*w.clusterName, *w.dataType, w.method).Inc()
221+
}
222+
223+
return nil
224+
}
225+
226+
// newWrappedCloudEventsMetricsStream creates a wrappedCloudEventsMetricsStream with the specified type and cluster reference.
227+
func newWrappedCloudEventsMetricsStream(clusterName, dataType *string, method string, ctx context.Context, ss grpc.ServerStream) grpc.ServerStream {
228+
return &wrappedCloudEventsMetricsStream{clusterName, dataType, method, ss, ctx}
229+
}
230+
231+
// NewCloudEventsMetricsStreamInterceptor creates a stream server interceptor for server metrics.
232+
func NewCloudEventsMetricsStreamInterceptor() grpc.StreamServerInterceptor {
233+
return func(srv interface{}, stream grpc.ServerStream, info *grpc.StreamServerInfo, handler grpc.StreamHandler) error {
234+
// extract the service and method
235+
service, method := SplitMethod(info.FullMethod)
236+
if service != gRPCCloudEventService {
237+
// skip emitting cloudevents metrics if it's not a cloudevents RPC request
238+
return handler(srv, stream)
239+
}
240+
241+
dataType := ""
242+
cluster := ""
243+
// create a wrapped stream to capture the source and emit metrics
244+
wrappedCEMetricsStream := newWrappedCloudEventsMetricsStream(&cluster, &dataType, method, stream.Context(), stream)
245+
// call rpc handler to handle RPC request
246+
err := handler(srv, wrappedCEMetricsStream)
247+
248+
// get status code from error
249+
status := statusFromError(err)
250+
code := status.Code()
251+
grpcCEProcessedCountMetric.WithLabelValues(cluster, dataType, method, code.String()).Inc()
252+
253+
return err
254+
}
255+
}
256+
257+
// statusFromError returns a grpc status. If the error code is neither a valid grpc status
258+
// nor a context error, codes.Unknown will be set.
259+
func statusFromError(err error) *status.Status {
260+
s, ok := status.FromError(err)
261+
// mirror what the grpc server itself does, i.e. also convert context errors to status
262+
if !ok {
263+
s = status.FromContextError(err)
264+
}
265+
return s
266+
}
267+
268+
// SplitMethod parses grpc full method "/package.service/method" to service and method
269+
func SplitMethod(fullMethod string) (service, method string) {
270+
if fullMethod == "" {
271+
return "unknown", "unknown"
272+
}
273+
// remove leading "/"
274+
if fullMethod[0] == '/' {
275+
fullMethod = fullMethod[1:]
276+
}
277+
// split at last "/"
278+
for i := len(fullMethod) - 1; i >= 0; i-- {
279+
if fullMethod[i] == '/' {
280+
return fullMethod[:i], fullMethod[i+1:]
281+
}
282+
}
283+
284+
return fullMethod, "unknown"
285+
}
286+
287+
// Register all the grpc server metrics for cloudevents.
288+
func RegisterCloudEventsGRPCMetrics() {
289+
once.Do(func() {
290+
metrics := []k8smetrics.Registerable{
291+
grpcCECalledCountMetric,
292+
grpcCEProcessedCountMetric,
293+
grpcCEProcessingDurationMetric,
294+
grpcCEMessageReceivedCountMetric,
295+
grpcCEMessageSentCountMetric,
296+
}
297+
298+
for _, m := range metrics {
299+
if m != nil {
300+
legacyregistry.MustRegister(m)
301+
} else {
302+
klog.Errorf("failed to register nil grpc server metric")
303+
}
304+
}
305+
})
306+
}

pkg/server/grpc/metrics/stats_test.go renamed to pkg/cloudevents/server/grpc/metrics/metrics_test.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@ func TestSplitMethod(t *testing.T) {
1919

2020
for _, tt := range tests {
2121
t.Run(tt.name, func(t *testing.T) {
22-
gotService, gotMethod := splitMethod(tt.fullMethod)
22+
gotService, gotMethod := SplitMethod(tt.fullMethod)
2323
if gotService != tt.expectedService {
2424
t.Errorf("splitMethod(%s) gotService = %s, want %s", tt.fullMethod, gotService, tt.expectedService)
2525
}

0 commit comments

Comments
 (0)