diff --git a/pkg/server/distributor.go b/pkg/server/distributor.go index 4afc545227b..79306666822 100644 --- a/pkg/server/distributor.go +++ b/pkg/server/distributor.go @@ -6,26 +6,17 @@ import ( "github.com/grafana/dskit/services" "github.com/grafana/grafana/pkg/modules" "github.com/grafana/grafana/pkg/services/grpcserver" - "github.com/grafana/grafana/pkg/services/grpcserver/interceptors" "github.com/grafana/grafana/pkg/storage/unified/resource" - resourcegrpc "github.com/grafana/grafana/pkg/storage/unified/resource/grpc" - "github.com/grafana/grafana/pkg/storage/unified/sql" "go.opentelemetry.io/otel" ) func (ms *ModuleServer) initDistributor() (services.Service, error) { - distributor := &distributorService{} - - tracer := otel.Tracer("unified-storage-distributor") - // FIXME: This is a temporary solution while we are migrating to the new authn interceptor - // grpcutils.NewGrpcAuthenticator should be used instead. - authn := sql.NewAuthenticatorWithFallback(ms.cfg, ms.registerer, tracer, func(ctx context.Context) (context.Context, error) { - auth := resourcegrpc.Authenticator{Tracer: tracer} - return auth.Authenticate(ctx) - }) - - var err error - distributor.grpcHandler, err = resource.ProvideDistributorServer(ms.cfg, ms.features, interceptors.AuthenticatorFunc(authn), ms.registerer, tracer, ms.storageRing, ms.storageRingClientPool) + var ( + distributor = &distributorService{} + tracer = otel.Tracer("unified-storage-distributor") + err error + ) + distributor.grpcHandler, err = resource.ProvideDistributorServer(ms.cfg, ms.features, ms.registerer, tracer, ms.storageRing, ms.storageRingClientPool) if err != nil { return nil, err } diff --git a/pkg/server/ring.go b/pkg/server/ring.go index 1c5e4af860e..1499702a81a 100644 --- a/pkg/server/ring.go +++ b/pkg/server/ring.go @@ -124,10 +124,7 @@ func newClientPool(clientCfg grpcclient.Config, log log.Logger, reg prometheus.R return nil, fmt.Errorf("failed to dial resource server %s %s: %s", inst.Id, inst.Addr, err) } - client, err := resource.NewResourceClient(conn, cfg, features, tracer) - if err != nil { - return nil, fmt.Errorf("failed to get resource client for server %s %s: %s", inst.Id, inst.Addr, err) - } + client := resource.NewAuthlessResourceClient(conn) return &resource.RingClient{ Client: client, diff --git a/pkg/services/grpcserver/service.go b/pkg/services/grpcserver/service.go index 536bb2cdf8b..e42fb60b4a0 100644 --- a/pkg/services/grpcserver/service.go +++ b/pkg/services/grpcserver/service.go @@ -68,22 +68,28 @@ func ProvideService(cfg *setting.Cfg, features featuremgmt.FeatureToggles, authe } } + unaryInterceptors := []grpc.UnaryServerInterceptor{ + interceptors.LoggingUnaryInterceptor(s.logger, s.cfg.EnableLogging), // needs to be registered after tracing interceptor to get trace id + middleware.UnaryServerInstrumentInterceptor(grpcRequestDuration), + } + streamInterceptors := []grpc.StreamServerInterceptor{ + interceptors.TracingStreamInterceptor(tracer), + interceptors.LoggingStreamInterceptor(s.logger, s.cfg.EnableLogging), + middleware.StreamServerInstrumentInterceptor(grpcRequestDuration), + } + + if authenticator != nil { + unaryInterceptors = append([]grpc.UnaryServerInterceptor{grpcAuth.UnaryServerInterceptor(authenticator.Authenticate)}, unaryInterceptors...) + streamInterceptors = append([]grpc.StreamServerInterceptor{grpcAuth.StreamServerInterceptor(authenticator.Authenticate)}, streamInterceptors...) + } + // Default auth is admin token check, but this can be overridden by // services which implement ServiceAuthFuncOverride interface. // See https://github.com/grpc-ecosystem/go-grpc-middleware/blob/main/interceptors/auth/auth.go#L30. opts := []grpc.ServerOption{ grpc.StatsHandler(otelgrpc.NewServerHandler()), - grpc.ChainUnaryInterceptor( - grpcAuth.UnaryServerInterceptor(authenticator.Authenticate), - interceptors.LoggingUnaryInterceptor(s.logger, s.cfg.EnableLogging), // needs to be registered after tracing interceptor to get trace id - middleware.UnaryServerInstrumentInterceptor(grpcRequestDuration), - ), - grpc.ChainStreamInterceptor( - interceptors.TracingStreamInterceptor(tracer), - interceptors.LoggingStreamInterceptor(s.logger, s.cfg.EnableLogging), - grpcAuth.StreamServerInterceptor(authenticator.Authenticate), - middleware.StreamServerInstrumentInterceptor(grpcRequestDuration), - ), + grpc.ChainUnaryInterceptor(unaryInterceptors...), + grpc.ChainStreamInterceptor(streamInterceptors...), } if s.cfg.TLSConfig != nil { diff --git a/pkg/storage/unified/resource/client.go b/pkg/storage/unified/resource/client.go index 8f71a0a4fe0..37ca1d2e500 100644 --- a/pkg/storage/unified/resource/client.go +++ b/pkg/storage/unified/resource/client.go @@ -64,8 +64,7 @@ func NewResourceClient(conn grpc.ClientConnInterface, cfg *setting.Cfg, features }) } -func NewLegacyResourceClient(channel grpc.ClientConnInterface) ResourceClient { - cc := grpchan.InterceptClientConn(channel, grpcUtils.UnaryClientInterceptor, grpcUtils.StreamClientInterceptor) +func newResourceClient(cc grpc.ClientConnInterface) ResourceClient { return &resourceClient{ ResourceStoreClient: resourcepb.NewResourceStoreClient(cc), ResourceIndexClient: resourcepb.NewResourceIndexClient(cc), @@ -76,6 +75,15 @@ func NewLegacyResourceClient(channel grpc.ClientConnInterface) ResourceClient { } } +func NewAuthlessResourceClient(cc grpc.ClientConnInterface) ResourceClient { + return newResourceClient(cc) +} + +func NewLegacyResourceClient(channel grpc.ClientConnInterface) ResourceClient { + cc := grpchan.InterceptClientConn(channel, grpcUtils.UnaryClientInterceptor, grpcUtils.StreamClientInterceptor) + return newResourceClient(cc) +} + func NewLocalResourceClient(server ResourceServer) ResourceClient { // scenario: local in-proc channel := &inprocgrpc.Channel{} @@ -106,14 +114,7 @@ func NewLocalResourceClient(server ResourceServer) ResourceClient { ) cc := grpchan.InterceptClientConn(channel, clientInt.UnaryClientInterceptor, clientInt.StreamClientInterceptor) - return &resourceClient{ - ResourceStoreClient: resourcepb.NewResourceStoreClient(cc), - ResourceIndexClient: resourcepb.NewResourceIndexClient(cc), - ManagedObjectIndexClient: resourcepb.NewManagedObjectIndexClient(cc), - BulkStoreClient: resourcepb.NewBulkStoreClient(cc), - BlobStoreClient: resourcepb.NewBlobStoreClient(cc), - DiagnosticsClient: resourcepb.NewDiagnosticsClient(cc), - } + return newResourceClient(cc) } type RemoteResourceClientConfig struct { diff --git a/pkg/storage/unified/resource/distributor.go b/pkg/storage/unified/resource/distributor.go index 7588a1f076f..231a0431597 100644 --- a/pkg/storage/unified/resource/distributor.go +++ b/pkg/storage/unified/resource/distributor.go @@ -12,18 +12,18 @@ import ( "github.com/grafana/grafana/pkg/infra/log" "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/services/grpcserver" - "github.com/grafana/grafana/pkg/services/grpcserver/interceptors" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/unified/resourcepb" "github.com/prometheus/client_golang/prometheus" "go.opentelemetry.io/otel/trace" "google.golang.org/grpc" "google.golang.org/grpc/health/grpc_health_v1" + "google.golang.org/grpc/metadata" ) -func ProvideDistributorServer(cfg *setting.Cfg, features featuremgmt.FeatureToggles, authnInterceptor interceptors.Authenticator, registerer prometheus.Registerer, tracer trace.Tracer, ring *ring.Ring, ringClientPool *ringclient.Pool) (grpcserver.Provider, error) { +func ProvideDistributorServer(cfg *setting.Cfg, features featuremgmt.FeatureToggles, registerer prometheus.Registerer, tracer trace.Tracer, ring *ring.Ring, ringClientPool *ringclient.Pool) (grpcserver.Provider, error) { var err error - grpcHandler, err := grpcserver.ProvideService(cfg, features, authnInterceptor, tracer, registerer) + grpcHandler, err := grpcserver.ProvideService(cfg, features, nil, tracer, registerer) if err != nil { return nil, err } @@ -239,9 +239,14 @@ func (ds *distributorServer) getClientToDistributeRequest(ctx context.Context, n return ctx, nil, err } + md, ok := metadata.FromIncomingContext(ctx) + if !ok { + md = make(metadata.MD) + } + ds.log.Info("distributing request to ", "methodName", methodName, "instanceId", rs.Instances[0].Id) - return userutils.InjectOrgID(ctx, namespace), client.(*RingClient).Client, nil + return userutils.InjectOrgID(metadata.NewOutgoingContext(ctx, md), namespace), client.(*RingClient).Client, nil } func (ds *distributorServer) IsHealthy(ctx context.Context, r *resourcepb.HealthCheckRequest) (*resourcepb.HealthCheckResponse, error) {