From 156acd17444c586b62e9279da5a7a5175f02c819 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Peter=20=C5=A0tibran=C3=BD?= Date: Thu, 19 Jun 2025 16:38:12 +0200 Subject: [PATCH] Add tracing to live calls. (#106980) * Add tracing to live calls. Clients using /api/live/ws use long-running request and post various commands via the request. This PR adds tracing span to each client-initiated action like subscribing/unsubscribing to/from channel, RPC call, publishing an event. Server-initiated messages are not included in the trace yet. --- pkg/services/live/live.go | 127 +++++++++++++++++++++++++++++++------- 1 file changed, 104 insertions(+), 23 deletions(-) diff --git a/pkg/services/live/live.go b/pkg/services/live/live.go index 8d101749527..116726a9e86 100644 --- a/pkg/services/live/live.go +++ b/pkg/services/live/live.go @@ -17,6 +17,10 @@ import ( "github.com/go-redis/redis/v8" "github.com/gobwas/glob" jsoniter "github.com/json-iterator/go" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/attribute" + "go.opentelemetry.io/otel/codes" + "go.opentelemetry.io/otel/trace" "golang.org/x/sync/errgroup" "github.com/grafana/grafana-plugin-sdk-go/backend" @@ -65,6 +69,7 @@ import ( var ( logger = log.New("live") loggerCF = log.New("live.centrifuge") + tracer = otel.Tracer("github.com/grafana/grafana/pkg/services/live") ) // CoreGrafanaScope list of core features @@ -208,12 +213,20 @@ func ProvideService(plugCtxProvider *plugincontext.Provider, cfg *setting.Cfg, r // different goroutines (belonging to different client connections). This is also // true for other event handlers. node.OnConnect(func(client *centrifuge.Client) { + _, connectSpan := tracer.Start(client.Context(), "live.OnConnect") + defer connectSpan.End() + connectSpan.SetAttributes( + attribute.String("user", client.UserID()), + attribute.String("client", client.ID()), + ) + numConnections := g.node.Hub().NumClients() if g.Cfg.LiveMaxConnections >= 0 && numConnections > g.Cfg.LiveMaxConnections { logger.Warn( "Max number of Live connections reached, increase max_connections in [live] configuration section", "user", client.UserID(), "client", client.ID(), "limit", g.Cfg.LiveMaxConnections, ) + connectSpan.AddEvent("disconnect", trace.WithAttributes(attribute.String("reason", "connection limit reached"))) client.Disconnect(centrifuge.DisconnectConnectionLimit) return } @@ -226,21 +239,58 @@ func ProvideService(plugCtxProvider *plugincontext.Provider, cfg *setting.Cfg, r // Called when client issues RPC (async request over Live connection). client.OnRPC(func(e centrifuge.RPCEvent, cb centrifuge.RPCCallback) { - err := runConcurrentlyIfNeeded(client.Context(), semaphore, func() { - cb(g.handleOnRPC(client, e)) + ctx, span := tracer.Start(client.Context(), "live.OnRPC") + // We finish span when calling callback, which can be done on a separate goroutine. + + span.SetAttributes( + attribute.String("method", e.Method), + attribute.String("data", string(e.Data)), + ) + + cbWithSpan := func(resp centrifuge.RPCReply, err error) { + defer span.End() + if err != nil { + span.SetStatus(codes.Error, err.Error()) + } else { + span.AddEvent("result", trace.WithAttributes(attribute.String("data", string(resp.Data)))) + span.SetStatus(codes.Ok, "") + } + cb(resp, err) + } + + err := runConcurrentlyIfNeeded(ctx, semaphore, func() { + cbWithSpan(g.handleOnRPC(ctx, client, e)) }) if err != nil { - cb(centrifuge.RPCReply{}, err) + cbWithSpan(centrifuge.RPCReply{}, err) } }) // Called when client subscribes to the channel. client.OnSubscribe(func(e centrifuge.SubscribeEvent, cb centrifuge.SubscribeCallback) { - err := runConcurrentlyIfNeeded(client.Context(), semaphore, func() { - cb(g.handleOnSubscribe(context.Background(), client, e)) + ctx, span := tracer.Start(client.Context(), "live.OnSubscribe") + // We finish span when calling callback, which can be done on a separate goroutine. + + span.SetAttributes( + attribute.String("channel", e.Channel), + attribute.String("data", string(e.Data)), + ) + + cbWithSpan := func(resp centrifuge.SubscribeReply, err error) { + defer span.End() + if err != nil { + span.SetStatus(codes.Error, err.Error()) + } else { + span.SetStatus(codes.Ok, "") + } + cb(resp, err) + } + + err := runConcurrentlyIfNeeded(ctx, semaphore, func() { + cbWithSpan(g.handleOnSubscribe(ctx, client, e)) }) if err != nil { - cb(centrifuge.SubscribeReply{}, err) + cbWithSpan(centrifuge.SubscribeReply{}, err) } }) @@ -248,15 +298,46 @@ func ProvideService(plugCtxProvider *plugincontext.Provider, cfg *setting.Cfg, r // In general, we should prefer writing to the HTTP API, but this // allows some simple prototypes to work quickly. client.OnPublish(func(e centrifuge.PublishEvent, cb centrifuge.PublishCallback) { - err := runConcurrentlyIfNeeded(client.Context(), semaphore, func() { - cb(g.handleOnPublish(context.Background(), client, e)) + ctx, span := tracer.Start(client.Context(), "live.OnPublish") + // We finish span when calling callback, which can be done on a separate goroutine. + + span.SetAttributes( + attribute.String("channel", e.Channel), + attribute.String("data", string(e.Data)), + ) + + cbWithSpan := func(resp centrifuge.PublishReply, err error) { + defer span.End() + if err != nil { + span.SetStatus(codes.Error, err.Error()) + } else { + span.SetStatus(codes.Ok, "") + } + cb(resp, err) + } + + err := runConcurrentlyIfNeeded(ctx, semaphore, func() { + cbWithSpan(g.handleOnPublish(ctx, client, e)) }) if err != nil { - cb(centrifuge.PublishReply{}, err) + cbWithSpan(centrifuge.PublishReply{}, err) } }) + // We don't need to do anything on unsubscribe, but we create tracing span with channel name. + client.OnUnsubscribe(func(e centrifuge.UnsubscribeEvent) { + _, span := tracer.Start(client.Context(), "live.OnUnsubscribe") + defer span.End() + + span.SetAttributes( + attribute.String("channel", e.Channel), + ) + }) + client.OnDisconnect(func(e centrifuge.DisconnectEvent) { + _, span := tracer.Start(client.Context(), "live.OnDisconnect") + defer span.End() + reason := e.Reason if e.Code == 3001 { // Shutdown return @@ -585,12 +666,12 @@ func (g *GrafanaLive) HandleDatasourceUpdate(orgID int64, dsUID string) { // that map keys is ordered. var jsonStd = jsoniter.ConfigCompatibleWithStandardLibrary -func (g *GrafanaLive) handleOnRPC(client *centrifuge.Client, e centrifuge.RPCEvent) (centrifuge.RPCReply, error) { +func (g *GrafanaLive) handleOnRPC(clientContextWithSpan context.Context, client *centrifuge.Client, e centrifuge.RPCEvent) (centrifuge.RPCReply, error) { logger.Debug("Client calls RPC", "user", client.UserID(), "client", client.ID(), "method", e.Method) if e.Method != "grafana.query" { return centrifuge.RPCReply{}, centrifuge.ErrorMethodNotFound } - user, ok := livecontext.GetContextSignedUser(client.Context()) + user, ok := livecontext.GetContextSignedUser(clientContextWithSpan) if !ok { logger.Error("No user found in context", "user", client.UserID(), "client", client.ID(), "method", e.Method) return centrifuge.RPCReply{}, centrifuge.ErrorInternal @@ -600,7 +681,7 @@ func (g *GrafanaLive) handleOnRPC(client *centrifuge.Client, e centrifuge.RPCEve if err != nil { return centrifuge.RPCReply{}, centrifuge.ErrorBadRequest } - resp, err := g.queryDataService.QueryData(client.Context(), user, false, req) + resp, err := g.queryDataService.QueryData(clientContextWithSpan, user, false, req) if err != nil { logger.Error("Error query data", "user", client.UserID(), "client", client.ID(), "method", e.Method, "error", err) if errors.Is(err, datasources.ErrDataSourceAccessDenied) { @@ -622,10 +703,10 @@ func (g *GrafanaLive) handleOnRPC(client *centrifuge.Client, e centrifuge.RPCEve }, nil } -func (g *GrafanaLive) handleOnSubscribe(ctx context.Context, client *centrifuge.Client, e centrifuge.SubscribeEvent) (centrifuge.SubscribeReply, error) { +func (g *GrafanaLive) handleOnSubscribe(clientContextWithSpan context.Context, client *centrifuge.Client, e centrifuge.SubscribeEvent) (centrifuge.SubscribeReply, error) { logger.Debug("Client wants to subscribe", "user", client.UserID(), "client", client.ID(), "channel", e.Channel) - user, ok := livecontext.GetContextSignedUser(client.Context()) + user, ok := livecontext.GetContextSignedUser(clientContextWithSpan) if !ok { logger.Error("No user found in context", "user", client.UserID(), "client", client.ID(), "channel", e.Channel) return centrifuge.SubscribeReply{}, centrifuge.ErrorInternal @@ -656,7 +737,7 @@ func (g *GrafanaLive) handleOnSubscribe(ctx context.Context, client *centrifuge. ruleFound = ok if ok { if rule.SubscribeAuth != nil { - ok, err := rule.SubscribeAuth.CanSubscribe(client.Context(), user) + ok, err := rule.SubscribeAuth.CanSubscribe(clientContextWithSpan, user) if err != nil { logger.Error("Error checking subscribe permissions", "user", client.UserID(), "client", client.ID(), "channel", e.Channel, "error", err) return centrifuge.SubscribeReply{}, centrifuge.ErrorInternal @@ -670,7 +751,7 @@ func (g *GrafanaLive) handleOnSubscribe(ctx context.Context, client *centrifuge. if len(rule.Subscribers) > 0 { var err error for _, sub := range rule.Subscribers { - reply, status, err = sub.Subscribe(client.Context(), pipeline.Vars{ + reply, status, err = sub.Subscribe(clientContextWithSpan, pipeline.Vars{ OrgID: orgID, Channel: channel, }, e.Data) @@ -686,7 +767,7 @@ func (g *GrafanaLive) handleOnSubscribe(ctx context.Context, client *centrifuge. } } if !ruleFound { - handler, addr, err := g.GetChannelHandler(ctx, user, channel) + handler, addr, err := g.GetChannelHandler(clientContextWithSpan, user, channel) if err != nil { if errors.Is(err, live.ErrInvalidChannelID) { logger.Info("Invalid channel ID", "user", client.UserID(), "client", client.ID(), "channel", e.Channel) @@ -695,7 +776,7 @@ func (g *GrafanaLive) handleOnSubscribe(ctx context.Context, client *centrifuge. logger.Error("Error getting channel handler", "user", client.UserID(), "client", client.ID(), "channel", e.Channel, "error", err) return centrifuge.SubscribeReply{}, centrifuge.ErrorInternal } - reply, status, err = handler.OnSubscribe(client.Context(), user, model.SubscribeEvent{ + reply, status, err = handler.OnSubscribe(clientContextWithSpan, user, model.SubscribeEvent{ Channel: channel, Path: addr.Path, Data: e.Data, @@ -723,10 +804,10 @@ func (g *GrafanaLive) handleOnSubscribe(ctx context.Context, client *centrifuge. }, nil } -func (g *GrafanaLive) handleOnPublish(ctx context.Context, client *centrifuge.Client, e centrifuge.PublishEvent) (centrifuge.PublishReply, error) { +func (g *GrafanaLive) handleOnPublish(clientCtxWithSpan context.Context, client *centrifuge.Client, e centrifuge.PublishEvent) (centrifuge.PublishReply, error) { logger.Debug("Client wants to publish", "user", client.UserID(), "client", client.ID(), "channel", e.Channel) - user, ok := livecontext.GetContextSignedUser(client.Context()) + user, ok := livecontext.GetContextSignedUser(clientCtxWithSpan) if !ok { logger.Error("No user found in context", "user", client.UserID(), "client", client.ID(), "channel", e.Channel) return centrifuge.PublishReply{}, centrifuge.ErrorInternal @@ -752,7 +833,7 @@ func (g *GrafanaLive) handleOnPublish(ctx context.Context, client *centrifuge.Cl } if ok { if rule.PublishAuth != nil { - ok, err := rule.PublishAuth.CanPublish(client.Context(), user) + ok, err := rule.PublishAuth.CanPublish(clientCtxWithSpan, user) if err != nil { logger.Error("Error checking publish permissions", "user", client.UserID(), "client", client.ID(), "channel", e.Channel, "error", err) return centrifuge.PublishReply{}, centrifuge.ErrorInternal @@ -769,7 +850,7 @@ func (g *GrafanaLive) handleOnPublish(ctx context.Context, client *centrifuge.Cl return centrifuge.PublishReply{}, ¢rifuge.Error{Code: uint32(code), Message: text} } } - _, err := g.Pipeline.ProcessInput(client.Context(), user.GetOrgID(), channel, e.Data) + _, err := g.Pipeline.ProcessInput(clientCtxWithSpan, user.GetOrgID(), channel, e.Data) if err != nil { logger.Error("Error processing input", "user", client.UserID(), "client", client.ID(), "channel", e.Channel, "error", err) return centrifuge.PublishReply{}, centrifuge.ErrorInternal @@ -780,7 +861,7 @@ func (g *GrafanaLive) handleOnPublish(ctx context.Context, client *centrifuge.Cl } } - handler, addr, err := g.GetChannelHandler(ctx, user, channel) + handler, addr, err := g.GetChannelHandler(clientCtxWithSpan, user, channel) if err != nil { if errors.Is(err, live.ErrInvalidChannelID) { logger.Info("Invalid channel ID", "user", client.UserID(), "client", client.ID(), "channel", e.Channel)