Merge remote-tracking branch 'origin/main' into ds-apiserver-with-configs

This commit is contained in:
Ryan McKinley
2025-06-20 00:18:05 +03:00
152 changed files with 1064 additions and 976 deletions
+104 -23
View File
@@ -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{}, &centrifuge.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)