From c8da64e4ebc49a80483257846fc041065a881792 Mon Sep 17 00:00:00 2001 From: Rafael Paulovic Date: Thu, 15 Jan 2026 00:49:08 +0100 Subject: [PATCH] chore: separate resource service into search and storage --- pkg/server/search_server_distributor_test.go | 27 +- pkg/server/wire_gen.go | 23 +- pkg/server/wireexts_oss.go | 15 +- pkg/setting/setting.go | 2 - pkg/setting/setting_unified_storage.go | 4 - pkg/storage/unified/client.go | 66 +++-- pkg/storage/unified/resource/server.go | 3 - pkg/storage/unified/sql/backend.go | 62 ++++- pkg/storage/unified/sql/remote_search.go | 90 ------- pkg/storage/unified/sql/search_service.go | 169 ++++++------- pkg/storage/unified/sql/server.go | 91 ++----- pkg/storage/unified/sql/service.go | 196 ++++++++------- pkg/storage/unified/sql/storage_service.go | 234 ++++++++++++++++++ .../unified/testing/search_and_storage.go | 3 +- 14 files changed, 596 insertions(+), 389 deletions(-) delete mode 100644 pkg/storage/unified/sql/remote_search.go create mode 100644 pkg/storage/unified/sql/storage_service.go diff --git a/pkg/server/search_server_distributor_test.go b/pkg/server/search_server_distributor_test.go index 7c572256383..09add10350f 100644 --- a/pkg/server/search_server_distributor_test.go +++ b/pkg/server/search_server_distributor_test.go @@ -391,17 +391,20 @@ func createBaselineServer(t *testing.T, dbType, dbConnStr string, testNamespaces require.NoError(t, err) searchOpts, err := search.NewSearchOptions(features, cfg, docBuilders, nil, nil) require.NoError(t, err) - server, err := sql.NewResourceServer(sql.ServerOptions{ - DB: nil, - Cfg: cfg, - Tracer: tracer, - Reg: nil, - AccessClient: nil, - SearchOptions: searchOpts, - StorageMetrics: nil, - IndexMetrics: nil, - Features: features, - QOSQueue: nil, + searchServer, err := sql.NewSearchServer(sql.SearchServerOptions{ + DB: nil, + Cfg: cfg, + Tracer: tracer, + Reg: nil, + AccessClient: nil, + SearchOptions: searchOpts, + IndexMetrics: nil, + }) + require.NoError(t, err) + storageServer, err := sql.NewStorageServer(sql.StorageServerOptions{ + Cfg: cfg, + Tracer: tracer, + Features: features, }) require.NoError(t, err) @@ -417,7 +420,7 @@ func createBaselineServer(t *testing.T, dbType, dbConnStr string, testNamespaces for _, ns := range testNamespaces { for range rand.Intn(maxPlaylistPerNamespace) + 1 { - _, err = server.Create(ctx, generatePlaylistPayload(ns)) + _, err = storageServer.Create(ctx, generatePlaylistPayload(ns)) require.NoError(t, err) } } diff --git a/pkg/server/wire_gen.go b/pkg/server/wire_gen.go index 65e6ebeb36e..c9625abf743 100644 --- a/pkg/server/wire_gen.go +++ b/pkg/server/wire_gen.go @@ -513,6 +513,11 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api if err != nil { return nil, err } + storageMetrics := resource.ProvideStorageMetrics(registerer) + storageBackend, err := sql.ProvideStorageBackend(cfg, sqlStore, tracer, registerer, storageMetrics) + if err != nil { + return nil, err + } options := &unified.Options{ Cfg: cfg, Features: featureToggles, @@ -522,8 +527,8 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api Authzc: accessClient, Docs: documentBuilderSupplier, SecureValues: inlineSecureValueSupport, + Backend: storageBackend, } - storageMetrics := resource.ProvideStorageMetrics(registerer) bleveIndexMetrics := resource.ProvideIndexMetrics(registerer) resourceClient, err := unified.ProvideUnifiedStorageClient(options, storageMetrics, bleveIndexMetrics) if err != nil { @@ -1173,6 +1178,11 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac if err != nil { return nil, err } + storageMetrics := resource.ProvideStorageMetrics(registerer) + storageBackend, err := sql.ProvideStorageBackend(cfg, sqlStore, tracer, registerer, storageMetrics) + if err != nil { + return nil, err + } options := &unified.Options{ Cfg: cfg, Features: featureToggles, @@ -1182,8 +1192,8 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac Authzc: accessClient, Docs: documentBuilderSupplier, SecureValues: inlineSecureValueSupport, + Backend: storageBackend, } - storageMetrics := resource.ProvideStorageMetrics(registerer) bleveIndexMetrics := resource.ProvideIndexMetrics(registerer) resourceClient, err := unified.ProvideUnifiedStorageClient(options, storageMetrics, bleveIndexMetrics) if err != nil { @@ -1748,7 +1758,14 @@ func InitializeModuleServer(cfg *setting.Cfg, opts Options, apiOpts api.ServerOp hooksService := hooks.ProvideService() ossLicensingService := licensing.ProvideService(cfg, hooksService) moduleRegisterer := ProvideNoopModuleRegisterer() - storageBackend, err := sql.ProvideStorageBackend(cfg) + ossMigrations := migrations.ProvideOSSMigrations(featureToggles) + inProcBus := bus.ProvideBus(tracingService) + sqlStore, err := sqlstore.ProvideService(cfg, featureToggles, ossMigrations, inProcBus, tracingService) + if err != nil { + return nil, err + } + tracer := otelTracer() + storageBackend, err := sql.ProvideStorageBackend(cfg, sqlStore, tracer, registerer, storageMetrics) if err != nil { return nil, err } diff --git a/pkg/server/wireexts_oss.go b/pkg/server/wireexts_oss.go index 4a12bf8db00..0cdb7de463f 100644 --- a/pkg/server/wireexts_oss.go +++ b/pkg/server/wireexts_oss.go @@ -6,6 +6,9 @@ package server import ( "github.com/google/wire" + "github.com/grafana/grafana/pkg/bus" + "github.com/grafana/grafana/pkg/infra/db" + "github.com/grafana/grafana/pkg/services/sqlstore" "github.com/grafana/grafana/pkg/configprovider" "github.com/grafana/grafana/pkg/infra/metrics" @@ -146,10 +149,10 @@ var wireExtsBasicSet = wire.NewSet( sandbox.ProvideService, wire.Bind(new(sandbox.Sandbox), new(*sandbox.Service)), wire.Struct(new(unified.Options), "*"), + sql.ProvideStorageBackend, unified.ProvideUnifiedStorageClient, wire.Bind(new(resourcepb.ResourceIndexClient), new(resource.ResourceClient)), wire.Bind(new(resource.MigratorClient), new(resource.ResourceClient)), - sql.ProvideStorageBackend, builder.ProvideDefaultBuildHandlerChainFuncFromBuilders, aggregatorrunner.ProvideNoopAggregatorConfigurator, apisregistry.WireSetExts, @@ -198,6 +201,16 @@ var wireExtsModuleServerSet = wire.NewSet( tracing.ProvideTracingConfig, tracing.ProvideService, wire.Bind(new(tracing.Tracer), new(*tracing.TracingService)), + otelTracer, + // Bus + bus.ProvideBus, + wire.Bind(new(bus.Bus), new(*bus.InProcBus)), + // Database migrations + migrations.ProvideOSSMigrations, + wire.Bind(new(registry.DatabaseMigrator), new(*migrations.OSSMigrations)), + // Database + sqlstore.ProvideService, + wire.Bind(new(db.DB), new(*sqlstore.SQLStore)), // Unified storage resource.ProvideStorageMetrics, resource.ProvideIndexMetrics, diff --git a/pkg/setting/setting.go b/pkg/setting/setting.go index b16f3a3e2ca..9667b82b9fa 100644 --- a/pkg/setting/setting.go +++ b/pkg/setting/setting.go @@ -623,8 +623,6 @@ type Cfg struct { OverridesFilePath string OverridesReloadInterval time.Duration EnableSQLKVBackend bool - SearchMode string // "", "embedded", "remote" - empty defaults to embedded - SearchServerAddress string // gRPC address for remote search server // Secrets Management SecretsManagement SecretsManagerSettings diff --git a/pkg/setting/setting_unified_storage.go b/pkg/setting/setting_unified_storage.go index fe1f021a6c0..21a3f455993 100644 --- a/pkg/setting/setting_unified_storage.go +++ b/pkg/setting/setting_unified_storage.go @@ -136,10 +136,6 @@ func (cfg *Cfg) setUnifiedStorageConfig() { // use sqlkv (resource/sqlkv) instead of the sql backend (sql/backend) as the StorageServer cfg.EnableSQLKVBackend = section.Key("enable_sqlkv_backend").MustBool(false) - // search mode: "", "embedded", "remote" - empty defaults to embedded for backward compatibility - cfg.SearchMode = section.Key("search_mode").MustString("") - cfg.SearchServerAddress = section.Key("search_server_address").String() - cfg.MaxFileIndexAge = section.Key("max_file_index_age").MustDuration(0) cfg.MinFileIndexBuildVersion = section.Key("min_file_index_build_version").MustString("") } diff --git a/pkg/storage/unified/client.go b/pkg/storage/unified/client.go index edbcd434210..909d322e6a9 100644 --- a/pkg/storage/unified/client.go +++ b/pkg/storage/unified/client.go @@ -46,6 +46,7 @@ type Options struct { Authzc types.AccessClient Docs resource.DocumentBuilderSupplier SecureValues secrets.InlineSecureValueSupport + Backend resource.StorageBackend // Shared backend to avoid duplicate metrics registration } type clientMetrics struct { @@ -67,7 +68,7 @@ func ProvideUnifiedStorageClient(opts *Options, BlobStoreURL: apiserverCfg.Key("blob_url").MustString(""), BlobThresholdBytes: apiserverCfg.Key("blob_threshold_bytes").MustInt(options.BlobThresholdDefault), GrpcClientKeepaliveTime: apiserverCfg.Key("grpc_client_keepalive_time").MustDuration(0), - }, opts.Cfg, opts.Features, opts.DB, opts.Tracer, opts.Reg, opts.Authzc, opts.Docs, storageMetrics, indexMetrics, opts.SecureValues) + }, opts.Cfg, opts.Features, opts.DB, opts.Tracer, opts.Reg, opts.Authzc, opts.Docs, storageMetrics, indexMetrics, opts.SecureValues, opts.Backend) if err == nil { // Decide whether to disable SQL fallback stats per resource in Mode 5. // Otherwise we would still try to query the legacy SQL database in Mode 5. @@ -103,6 +104,7 @@ func newClient(opts options.StorageOptions, storageMetrics *resource.StorageMetrics, indexMetrics *resource.BleveIndexMetrics, secure secrets.InlineSecureValueSupport, + backend resource.StorageBackend, ) (resource.ResourceClient, error) { ctx := context.Background() @@ -169,24 +171,14 @@ func newClient(opts options.StorageOptions, return resource.NewResourceClient(conn, indexConn, cfg, features, tracer) default: + // Create search options for the search server searchOptions, err := search.NewSearchOptions(features, cfg, docs, indexMetrics, nil) if err != nil { return nil, err } - serverOptions := sql.ServerOptions{ - DB: db, - Cfg: cfg, - Tracer: tracer, - Reg: reg, - AccessClient: authzc, - SearchOptions: searchOptions, - StorageMetrics: storageMetrics, - IndexMetrics: indexMetrics, - Features: features, - SecureValues: secure, - } - + // Setup QOS queue if enabled + var qosQueue sql.QOSEnqueueDequeuer if cfg.QOSEnabled { qosReg := prometheus.WrapRegistererWithPrefix("resource_server_qos_", reg) queue := scheduler.NewQueue(&scheduler.QueueOptions{ @@ -197,7 +189,7 @@ func newClient(opts options.StorageOptions, if err := services.StartAndAwaitRunning(ctx, queue); err != nil { return nil, fmt.Errorf("failed to start queue: %w", err) } - scheduler, err := scheduler.NewScheduler(queue, &scheduler.Config{ + sched, err := scheduler.NewScheduler(queue, &scheduler.Config{ NumWorkers: cfg.QOSNumberWorker, Logger: cfg.Logger, }) @@ -205,31 +197,59 @@ func newClient(opts options.StorageOptions, return nil, fmt.Errorf("failed to create scheduler: %w", err) } - err = services.StartAndAwaitRunning(ctx, scheduler) + err = services.StartAndAwaitRunning(ctx, sched) if err != nil { return nil, fmt.Errorf("failed to start scheduler: %w", err) } - serverOptions.QOSQueue = queue + qosQueue = queue } - // only enable if an overrides file path is provided + // Setup overrides service if enabled + var overridesSvc *resource.OverridesService if cfg.OverridesFilePath != "" { - overridesSvc, err := resource.NewOverridesService(ctx, cfg.Logger, reg, tracer, resource.ReloadOptions{ + overridesSvc, err = resource.NewOverridesService(ctx, cfg.Logger, reg, tracer, resource.ReloadOptions{ FilePath: cfg.OverridesFilePath, ReloadPeriod: cfg.OverridesReloadInterval, }) if err != nil { return nil, err } - - serverOptions.OverridesService = overridesSvc } - server, searchServer, err := sql.NewResourceServer(serverOptions) + // Create the search server with shared backend + searchServer, err := sql.NewSearchServer(sql.SearchServerOptions{ + Backend: backend, // Use shared backend to avoid duplicate metrics registration + DB: db, + Cfg: cfg, + Tracer: tracer, + Reg: reg, + AccessClient: authzc, + SearchOptions: searchOptions, + IndexMetrics: indexMetrics, + }) if err != nil { return nil, err } - return resource.NewLocalResourceClient(server, searchServer), nil + + // Create the storage server with shared backend + storageServer, err := sql.NewStorageServer(sql.StorageServerOptions{ + Backend: backend, // Use shared backend to avoid duplicate metrics registration + DB: db, + Cfg: cfg, + Tracer: tracer, + Reg: reg, + AccessClient: authzc, + StorageMetrics: storageMetrics, + Features: features, + QOSQueue: qosQueue, + SecureValues: secure, + OverridesService: overridesSvc, + }) + if err != nil { + return nil, err + } + + return resource.NewLocalResourceClient(storageServer, searchServer), nil } } diff --git a/pkg/storage/unified/resource/server.go b/pkg/storage/unified/resource/server.go index 37e31a0492b..1c2424225e0 100644 --- a/pkg/storage/unified/resource/server.go +++ b/pkg/storage/unified/resource/server.go @@ -227,9 +227,6 @@ type ResourceServerOptions struct { // The blob configuration Blob BlobConfig - // Search options - Search SearchServer - // Quota service OverridesService *OverridesService diff --git a/pkg/storage/unified/sql/backend.go b/pkg/storage/unified/sql/backend.go index a129a01727a..ad0bdc2e8b8 100644 --- a/pkg/storage/unified/sql/backend.go +++ b/pkg/storage/unified/sql/backend.go @@ -11,6 +11,9 @@ import ( "time" "github.com/go-sql-driver/mysql" + infraDB "github.com/grafana/grafana/pkg/infra/db" + "github.com/grafana/grafana/pkg/services/sqlstore/migrator" + "github.com/grafana/grafana/pkg/storage/unified/sql/db/dbimpl" "github.com/jackc/pgx/v5/pgconn" "github.com/lib/pq" "github.com/prometheus/client_golang/prometheus" @@ -44,10 +47,63 @@ const defaultPrunerHistoryLimit = 20 func ProvideStorageBackend( cfg *setting.Cfg, + db infraDB.DB, + tracer trace.Tracer, + reg prometheus.Registerer, + storageMetrics *resource.StorageMetrics, ) (resource.StorageBackend, error) { - // TODO: make this the central place to provide SQL backend - // Currently it is skipped as we need to handle the cases of Diagnostics and Lifecycle - return nil, nil + // Create the resource DB + eDB, err := dbimpl.ProvideResourceDB(db, cfg, tracer) + if err != nil { + return nil, fmt.Errorf("failed to create resource DB: %w", err) + } + + // Check if HA is enabled + isHA := isHighAvailabilityEnabled( + cfg.SectionWithEnvOverrides("database"), + cfg.SectionWithEnvOverrides("resource_api"), + ) + + // Create the backend + backend, err := NewBackend(BackendOptions{ + DBProvider: eDB, + Reg: reg, + IsHA: isHA, + storageMetrics: storageMetrics, + LastImportTimeMaxAge: cfg.MaxFileIndexAge, + }) + if err != nil { + return nil, fmt.Errorf("failed to create backend: %w", err) + } + + // Initialize the backend + if err := backend.Init(context.Background()); err != nil { + return nil, fmt.Errorf("failed to initialize backend: %w", err) + } + + return backend, nil +} + +// isHighAvailabilityEnabled determines if high availability mode should +// be enabled based on database configuration. High availability is enabled +// by default except for SQLite databases. +func isHighAvailabilityEnabled(dbCfg, resourceAPICfg *setting.DynamicSection) bool { + // If the resource API is using a non-SQLite database, we assume it's in HA mode. + resourceDBType := resourceAPICfg.Key("db_type").String() + if resourceDBType != "" && resourceDBType != migrator.SQLite { + return true + } + + // Check in the config if HA is enabled - by default we always assume a HA setup. + isHA := dbCfg.Key("high_availability").MustBool(true) + + // SQLite is not possible to run in HA, so we force it to false. + databaseType := dbCfg.Key("type").String() + if databaseType == migrator.SQLite { + isHA = false + } + + return isHA } type Backend interface { diff --git a/pkg/storage/unified/sql/remote_search.go b/pkg/storage/unified/sql/remote_search.go deleted file mode 100644 index 46b7a19ac6e..00000000000 --- a/pkg/storage/unified/sql/remote_search.go +++ /dev/null @@ -1,90 +0,0 @@ -package sql - -import ( - "context" - "fmt" - - "google.golang.org/grpc" - "google.golang.org/grpc/credentials/insecure" - - "github.com/grafana/grafana/pkg/storage/unified/resource" - "github.com/grafana/grafana/pkg/storage/unified/resourcepb" -) - -var _ resource.SearchServer = (*remoteSearchClient)(nil) - -// remoteSearchClient wraps gRPC search clients to implement the SearchServer interface. -// This allows the storage server to delegate search operations to a remote search server. -type remoteSearchClient struct { - conn *grpc.ClientConn - index resourcepb.ResourceIndexClient - moiClient resourcepb.ManagedObjectIndexClient - diagnostics resourcepb.DiagnosticsClient -} - -// newRemoteSearchClient creates a new remote search client that connects to a search server at the given address. -func newRemoteSearchClient(address string) (*remoteSearchClient, error) { - if address == "" { - return nil, fmt.Errorf("search server address is required for remote search mode") - } - - conn, err := grpc.NewClient(address, - grpc.WithTransportCredentials(insecure.NewCredentials()), - grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`), - ) - if err != nil { - return nil, fmt.Errorf("failed to create gRPC connection to search server: %w", err) - } - - return &remoteSearchClient{ - conn: conn, - index: resourcepb.NewResourceIndexClient(conn), - moiClient: resourcepb.NewManagedObjectIndexClient(conn), - diagnostics: resourcepb.NewDiagnosticsClient(conn), - }, nil -} - -// Init implements resource.LifecycleHooks. -// For remote search, there's nothing to initialize locally. -func (r *remoteSearchClient) Init(ctx context.Context) error { - return nil -} - -// Stop implements resource.LifecycleHooks. -// Closes the gRPC connection. -func (r *remoteSearchClient) Stop(ctx context.Context) error { - if r.conn != nil { - return r.conn.Close() - } - return nil -} - -// Search implements resourcepb.ResourceIndexServer. -func (r *remoteSearchClient) Search(ctx context.Context, req *resourcepb.ResourceSearchRequest) (*resourcepb.ResourceSearchResponse, error) { - return r.index.Search(ctx, req) -} - -// GetStats implements resourcepb.ResourceIndexServer. -func (r *remoteSearchClient) GetStats(ctx context.Context, req *resourcepb.ResourceStatsRequest) (*resourcepb.ResourceStatsResponse, error) { - return r.index.GetStats(ctx, req) -} - -// RebuildIndexes implements resourcepb.ResourceIndexServer. -func (r *remoteSearchClient) RebuildIndexes(ctx context.Context, req *resourcepb.RebuildIndexesRequest) (*resourcepb.RebuildIndexesResponse, error) { - return r.index.RebuildIndexes(ctx, req) -} - -// CountManagedObjects implements resourcepb.ManagedObjectIndexServer. -func (r *remoteSearchClient) CountManagedObjects(ctx context.Context, req *resourcepb.CountManagedObjectsRequest) (*resourcepb.CountManagedObjectsResponse, error) { - return r.moiClient.CountManagedObjects(ctx, req) -} - -// ListManagedObjects implements resourcepb.ManagedObjectIndexServer. -func (r *remoteSearchClient) ListManagedObjects(ctx context.Context, req *resourcepb.ListManagedObjectsRequest) (*resourcepb.ListManagedObjectsResponse, error) { - return r.moiClient.ListManagedObjects(ctx, req) -} - -// IsHealthy implements resourcepb.DiagnosticsServer. -func (r *remoteSearchClient) IsHealthy(ctx context.Context, req *resourcepb.HealthCheckRequest) (*resourcepb.HealthCheckResponse, error) { - return r.diagnostics.IsHealthy(ctx, req) -} diff --git a/pkg/storage/unified/sql/search_service.go b/pkg/storage/unified/sql/search_service.go index cb196bdfa4d..b4fa1a81f43 100644 --- a/pkg/storage/unified/sql/search_service.go +++ b/pkg/storage/unified/sql/search_service.go @@ -5,18 +5,16 @@ import ( "errors" "fmt" "hash/fnv" - "net" - "os" - "strconv" - "time" + "net/http" + "github.com/gorilla/mux" + "github.com/grafana/grafana/pkg/storage/unified/search" "github.com/prometheus/client_golang/prometheus" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/trace" "google.golang.org/grpc/health/grpc_health_v1" "github.com/grafana/dskit/kv" - "github.com/grafana/dskit/netutil" "github.com/grafana/dskit/ring" "github.com/grafana/dskit/services" @@ -31,57 +29,44 @@ import ( "github.com/grafana/grafana/pkg/storage/unified/resource" "github.com/grafana/grafana/pkg/storage/unified/resource/grpc" "github.com/grafana/grafana/pkg/storage/unified/resourcepb" - "github.com/grafana/grafana/pkg/storage/unified/search" ) +var _ UnifiedStorageGrpcService = (*searchService)(nil) + +// operation used by the search-servers to check if they own the namespace var ( - _ UnifiedSearchGrpcService = (*searchService)(nil) + searchOwnerRead = ring.NewOp([]ring.InstanceState{ring.JOINING, ring.ACTIVE, ring.LEAVING}, nil) ) -// UnifiedSearchGrpcService is the interface for the standalone search gRPC service. -// This follows the same naming convention as UnifiedStorageGrpcService. -type UnifiedSearchGrpcService interface { - services.NamedService - - // GetAddress returns the address where this service is running - GetAddress() string -} - type searchService struct { *services.BasicService + backend resource.StorageBackend + cfg *setting.Cfg + features featuremgmt.FeatureToggles + db infraDB.DB + stopCh chan struct{} + stoppedCh chan error + handler grpcserver.Provider + tracing trace.Tracer + authenticator func(ctx context.Context) (context.Context, error) + httpServerRouter *mux.Router + + log log.Logger + reg prometheus.Registerer + docBuilders resource.DocumentBuilderSupplier + indexMetrics *resource.BleveIndexMetrics + searchRing *ring.Ring + + // Ring lifecycle and sharding support + ringLifecycler *ring.BasicLifecycler + // Subservices manager subservices *services.Manager subservicesWatcher *services.FailureWatcher hasSubservices bool - - cfg *setting.Cfg - features featuremgmt.FeatureToggles - db infraDB.DB - stopCh chan struct{} - stoppedCh chan error - - handler grpcserver.Provider - - tracing trace.Tracer - - authenticator func(ctx context.Context) (context.Context, error) - - log log.Logger - reg prometheus.Registerer - indexMetrics *resource.BleveIndexMetrics - - docBuilders resource.DocumentBuilderSupplier - - searchRing *ring.Ring - ringLifecycler *ring.BasicLifecycler - - backend resource.StorageBackend } -// ProvideUnifiedSearchGrpcService creates a standalone search gRPC service. -// This is used when running search-server as a separate target. -// It follows the same naming convention as ProvideUnifiedStorageGrpcService. func ProvideUnifiedSearchGrpcService( cfg *setting.Cfg, features featuremgmt.FeatureToggles, @@ -93,8 +78,10 @@ func ProvideUnifiedSearchGrpcService( searchRing *ring.Ring, memberlistKVConfig kv.Config, backend resource.StorageBackend, -) (UnifiedSearchGrpcService, error) { - tracer := otel.Tracer("search-server") + httpServerRouter *mux.Router, +) (UnifiedStorageGrpcService, error) { + var err error + tracer := otel.Tracer("unified-search-server") authn := NewAuthenticatorWithFallback(cfg, reg, tracer, func(ctx context.Context) (context.Context, error) { auth := grpc.Authenticator{Tracer: tracer} @@ -102,6 +89,7 @@ func ProvideUnifiedSearchGrpcService( }) s := &searchService{ + backend: backend, cfg: cfg, features: features, stopCh: make(chan struct{}), @@ -114,8 +102,8 @@ func ProvideUnifiedSearchGrpcService( docBuilders: docBuilders, indexMetrics: indexMetrics, searchRing: searchRing, + httpServerRouter: httpServerRouter, subservicesWatcher: services.NewFailureWatcher(), - backend: backend, } subservices := []services.Service{} @@ -130,11 +118,13 @@ func ProvideUnifiedSearchGrpcService( return nil, fmt.Errorf("failed to create KV store client: %s", err) } - lifecyclerCfg, err := toSearchLifecyclerConfig(cfg, log) + lifecyclerCfg, err := toLifecyclerConfig(cfg, log) if err != nil { return nil, fmt.Errorf("failed to initialize search-ring lifecycler config: %s", err) } + // Define lifecycler delegates in reverse order (last to be called defined first because they're + // chained via "next delegate"). delegate := ring.BasicLifecyclerDelegate(ring.NewInstanceRegisterDelegate(ring.JOINING, resource.RingNumTokens)) delegate = ring.NewLeaveOnStoppingDelegate(delegate, log) delegate = ring.NewAutoForgetDelegate(resource.RingHeartbeatTimeout*2, delegate, log) @@ -158,18 +148,36 @@ func ProvideUnifiedSearchGrpcService( if len(subservices) > 0 { s.hasSubservices = true - var err error s.subservices, err = services.NewManager(subservices...) if err != nil { return nil, fmt.Errorf("failed to create subservices manager: %w", err) } } + // This will be used when running as a dskit service s.BasicService = services.NewBasicService(s.starting, s.running, s.stopping).WithName(modules.SearchServer) + // Register HTTP endpoints if router is provided + s.RegisterHTTPEndpoints(httpServerRouter) + return s, nil } +func (s *searchService) PrepareDownscale(w http.ResponseWriter, r *http.Request) { + switch r.Method { + case http.MethodPost: + s.log.Info("Preparing for downscale. Will not keep instance in ring on shutdown.") + s.ringLifecycler.SetKeepInstanceInTheRingOnShutdown(false) + case http.MethodDelete: + s.log.Info("Downscale canceled. Will keep instance in ring on shutdown.") + s.ringLifecycler.SetKeepInstanceInTheRingOnShutdown(true) + case http.MethodGet: + // used for delayed downscale use case, which we don't support. Leaving here for completion sake + s.log.Info("Received GET request for prepare-downscale. Behavior not implemented.") + default: + } +} + func (s *searchService) OwnsIndex(key resource.NamespacedResource) (bool, error) { if s.searchRing == nil { return true, nil @@ -206,19 +214,26 @@ func (s *searchService) starting(ctx context.Context) error { return err } + // Create search options for the search server searchOptions, err := search.NewSearchOptions(s.features, s.cfg, s.docBuilders, s.indexMetrics, s.OwnsIndex) if err != nil { return err } - // Create search server - searchServer, err := resource.NewSearchServer(searchOptions, s.backend, authzClient, nil, s.indexMetrics, s.OwnsIndex) + // Create the search server + searchServer, err := NewSearchServer(SearchServerOptions{ + Backend: s.backend, + DB: s.db, + Cfg: s.cfg, + Tracer: s.tracing, + Reg: s.reg, + AccessClient: authzClient, + SearchOptions: searchOptions, + IndexMetrics: s.indexMetrics, + OwnsIndexFn: s.OwnsIndex, + }) if err != nil { - return fmt.Errorf("failed to create search server: %w", err) - } - - if err := searchServer.Init(ctx); err != nil { - return fmt.Errorf("failed to initialize search server: %w", err) + return err } s.handler, err = grpcserver.ProvideService(s.cfg, s.features, interceptors.AuthenticatorFunc(s.authenticator), s.tracing, prometheus.DefaultRegisterer) @@ -226,14 +241,16 @@ func (s *searchService) starting(ctx context.Context) error { return err } + healthService, err := resource.ProvideHealthService(searchServer) + if err != nil { + return err + } + srv := s.handler.GetServer() + // Register search services resourcepb.RegisterResourceIndexServer(srv, searchServer) resourcepb.RegisterManagedObjectIndexServer(srv, searchServer) resourcepb.RegisterDiagnosticsServer(srv, searchServer) - healthService, err := resource.ProvideHealthService(searchServer) - if err != nil { - return fmt.Errorf("failed to create health service: %w", err) - } grpc_health_v1.RegisterHealthServer(srv, healthService) // register reflection service @@ -269,6 +286,7 @@ func (s *searchService) starting(ctx context.Context) error { return nil } +// GetAddress returns the address of the gRPC server. func (s *searchService) GetAddress() string { return s.handler.GetAddress() } @@ -297,37 +315,8 @@ func (s *searchService) stopping(_ error) error { return nil } -func toSearchLifecyclerConfig(cfg *setting.Cfg, logger log.Logger) (ring.BasicLifecyclerConfig, error) { - instanceAddr, err := ring.GetInstanceAddr(cfg.MemberlistBindAddr, netutil.PrivateNetworkInterfacesWithFallback([]string{"eth0", "en0"}, logger), logger, true) - if err != nil { - return ring.BasicLifecyclerConfig{}, err +func (s *searchService) RegisterHTTPEndpoints(httpServerRouter *mux.Router) { + if httpServerRouter != nil && s.cfg.EnableSharding { + httpServerRouter.Path("/prepare-downscale").Methods("GET", "POST", "DELETE").Handler(http.HandlerFunc(s.PrepareDownscale)) } - - instanceId := cfg.InstanceID - if instanceId == "" { - hostname, err := os.Hostname() - if err != nil { - return ring.BasicLifecyclerConfig{}, err - } - instanceId = hostname - } - - _, grpcPortStr, err := net.SplitHostPort(cfg.GRPCServer.Address) - if err != nil { - return ring.BasicLifecyclerConfig{}, fmt.Errorf("could not get grpc port from grpc server address: %s", err) - } - - grpcPort, err := strconv.Atoi(grpcPortStr) - if err != nil { - return ring.BasicLifecyclerConfig{}, fmt.Errorf("error converting grpc address port to int: %s", err) - } - - return ring.BasicLifecyclerConfig{ - Addr: fmt.Sprintf("%s:%d", instanceAddr, grpcPort), - ID: instanceId, - HeartbeatPeriod: 15 * time.Second, - HeartbeatTimeout: resource.RingHeartbeatTimeout, - TokensObservePeriod: 0, - NumTokens: resource.RingNumTokens, - }, nil } diff --git a/pkg/storage/unified/sql/server.go b/pkg/storage/unified/sql/server.go index ccfcd4ec3ce..c2c721a9d32 100644 --- a/pkg/storage/unified/sql/server.go +++ b/pkg/storage/unified/sql/server.go @@ -16,7 +16,6 @@ import ( secrets "github.com/grafana/grafana/pkg/registry/apis/secret/contracts" inlinesecurevalue "github.com/grafana/grafana/pkg/registry/apis/secret/inline" "github.com/grafana/grafana/pkg/services/featuremgmt" - "github.com/grafana/grafana/pkg/services/sqlstore/migrator" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/unified/resource" "github.com/grafana/grafana/pkg/storage/unified/sql/db/dbimpl" @@ -43,8 +42,8 @@ type SearchServerOptions struct { OwnsIndexFn func(key resource.NamespacedResource) (bool, error) } -// ServerOptions contains the options for creating a new ResourceServer -type ServerOptions struct { +// StorageServerOptions contains the options for creating a storage-only server (without search) +type StorageServerOptions struct { Backend resource.StorageBackend OverridesService *resource.OverridesService DB infraDB.DB @@ -52,20 +51,19 @@ type ServerOptions struct { Tracer trace.Tracer Reg prometheus.Registerer AccessClient types.AccessClient - SearchOptions resource.SearchOptions StorageMetrics *resource.StorageMetrics - IndexMetrics *resource.BleveIndexMetrics Features featuremgmt.FeatureToggles QOSQueue QOSEnqueueDequeuer SecureValues secrets.InlineSecureValueSupport - OwnsIndexFn func(key resource.NamespacedResource) (bool, error) - // Search is an optional pre-created search server. If nil, one will be created. - Search resource.SearchServer } // NewSearchServer creates a new SearchServer with the given options. // This can be used to create a standalone search server or to create a search server // that will be passed to NewResourceServer. +// +// Important: When running in monolith mode, the backend should be provided by the caller +// to avoid duplicate metrics registration. Only in standalone microservice mode should +// this function create its own backend. func NewSearchServer(opts SearchServerOptions) (resource.SearchServer, error) { backend := opts.Backend if backend == nil { @@ -106,7 +104,13 @@ func NewSearchServer(opts SearchServerOptions) (resource.SearchServer, error) { return search, nil } -func NewResourceServer(opts ServerOptions) (resource.ResourceServer, resource.SearchServer, error) { +// NewStorageServer creates a storage-only server without search capabilities. +// Use this when you want to run storage and search as separate services. +// +// Important: When running in monolith mode, the backend should be provided by the caller +// to avoid duplicate metrics registration. Only in standalone microservice mode should +// this function create its own backend. +func NewStorageServer(opts StorageServerOptions) (resource.ResourceServer, error) { apiserverCfg := opts.Cfg.SectionWithEnvOverrides("grafana-apiserver") if opts.SecureValues == nil && opts.Cfg != nil && opts.Cfg.SecretsManagement.GrpcClientEnable { @@ -117,7 +121,7 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, resource.Se nil, // not needed for gRPC client mode ) if err != nil { - return nil, nil, fmt.Errorf("failed to create inline secure value service: %w", err) + return nil, fmt.Errorf("failed to create inline secure value service: %w", err) } opts.SecureValues = inlineSecureValueService } @@ -137,7 +141,7 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, resource.Se dir := strings.Replace(serverOptions.Blob.URL, "./data", opts.Cfg.DataPath, 1) err := os.MkdirAll(dir, 0700) if err != nil { - return nil, nil, err + return nil, err } serverOptions.Blob.URL = "file:///" + dir } @@ -150,17 +154,16 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, resource.Se if opts.Backend != nil { serverOptions.Backend = opts.Backend - // TODO: we should probably have a proper interface for diagnostics/lifecycle } else { eDB, err := dbimpl.ProvideResourceDB(opts.DB, opts.Cfg, opts.Tracer) if err != nil { - return nil, nil, err + return nil, err } if opts.Cfg.EnableSQLKVBackend { sqlkv, err := resource.NewSQLKV(eDB) if err != nil { - return nil, nil, fmt.Errorf("error creating sqlkv: %s", err) + return nil, fmt.Errorf("error creating sqlkv: %s", err) } kvBackendOpts := resource.KVBackendOptions{ @@ -172,12 +175,12 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, resource.Se ctx := context.Background() dbConn, err := eDB.Init(ctx) if err != nil { - return nil, nil, fmt.Errorf("error initializing DB: %w", err) + return nil, fmt.Errorf("error initializing DB: %w", err) } dialect := sqltemplate.DialectForDriver(dbConn.DriverName()) if dialect == nil { - return nil, nil, fmt.Errorf("unsupported database driver: %s", dbConn.DriverName()) + return nil, fmt.Errorf("unsupported database driver: %s", dbConn.DriverName()) } rvManager, err := rvmanager.NewResourceVersionManager(rvmanager.ResourceManagerOptions{ @@ -185,14 +188,13 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, resource.Se DB: dbConn, }) if err != nil { - return nil, nil, fmt.Errorf("failed to create resource version manager: %w", err) + return nil, fmt.Errorf("failed to create resource version manager: %w", err) } - // TODO add config to decide whether to pass RvManager or not kvBackendOpts.RvManager = rvManager kvBackend, err := resource.NewKVStorageBackend(kvBackendOpts) if err != nil { - return nil, nil, fmt.Errorf("error creating kv backend: %s", err) + return nil, fmt.Errorf("error creating kv backend: %s", err) } serverOptions.Backend = kvBackend @@ -206,10 +208,10 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, resource.Se Reg: opts.Reg, IsHA: isHA, storageMetrics: opts.StorageMetrics, - LastImportTimeMaxAge: opts.Cfg.MaxFileIndexAge, // No need to keep last_import_times older than max index age. + LastImportTimeMaxAge: opts.Cfg.MaxFileIndexAge, }) if err != nil { - return nil, nil, err + return nil, err } serverOptions.Backend = backend serverOptions.Diagnostics = backend @@ -217,56 +219,15 @@ func NewResourceServer(opts ServerOptions) (resource.ResourceServer, resource.Se } } - // Initialize the backend before creating search server (it needs the DB connection) + // Initialize the backend before creating server if serverOptions.Lifecycle != nil { if err := serverOptions.Lifecycle.Init(context.Background()); err != nil { - return nil, nil, fmt.Errorf("failed to initialize backend: %w", err) + return nil, fmt.Errorf("failed to initialize backend: %w", err) } } - // Use pre-created search server if provided, otherwise create one - search := opts.Search - if search == nil { - var err error - search, err = resource.NewSearchServer(opts.SearchOptions, serverOptions.Backend, opts.AccessClient, nil, opts.IndexMetrics, opts.OwnsIndexFn) - if err != nil { - return nil, nil, fmt.Errorf("failed to create search server: %w", err) - } - - if err := search.Init(context.Background()); err != nil { - return nil, nil, fmt.Errorf("failed to initialize search server: %w", err) - } - } - - serverOptions.Search = search serverOptions.QOSQueue = opts.QOSQueue serverOptions.OverridesService = opts.OverridesService - rs, err := resource.NewResourceServer(serverOptions) - if err != nil { - _ = search.Stop(context.Background()) - } - return rs, search, err -} - -// isHighAvailabilityEnabled determines if high availability mode should -// be enabled based on database configuration. High availability is enabled -// by default except for SQLite databases. -func isHighAvailabilityEnabled(dbCfg, resourceAPICfg *setting.DynamicSection) bool { - // If the resource API is using a non-SQLite database, we assume it's in HA mode. - resourceDBType := resourceAPICfg.Key("db_type").String() - if resourceDBType != "" && resourceDBType != migrator.SQLite { - return true - } - - // Check in the config if HA is enabled - by default we always assume a HA setup. - isHA := dbCfg.Key("high_availability").MustBool(true) - - // SQLite is not possible to run in HA, so we force it to false. - databaseType := dbCfg.Key("type").String() - if databaseType == migrator.SQLite { - isHA = false - } - - return isHA + return resource.NewResourceServer(serverOptions) } diff --git a/pkg/storage/unified/sql/service.go b/pkg/storage/unified/sql/service.go index 2d167a00a12..192cff25043 100644 --- a/pkg/storage/unified/sql/service.go +++ b/pkg/storage/unified/sql/service.go @@ -12,6 +12,7 @@ import ( "time" "github.com/gorilla/mux" + "github.com/grafana/grafana/pkg/storage/unified/search" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promauto" "go.opentelemetry.io/otel" @@ -36,6 +37,7 @@ import ( "github.com/grafana/grafana/pkg/storage/unified/resource" "github.com/grafana/grafana/pkg/storage/unified/resource/grpc" "github.com/grafana/grafana/pkg/storage/unified/resourcepb" + "github.com/grafana/grafana/pkg/storage/unified/sql/db/dbimpl" "github.com/grafana/grafana/pkg/util/scheduler" ) @@ -53,38 +55,40 @@ type UnifiedStorageGrpcService interface { type service struct { *services.BasicService - // Subservices manager - subservices *services.Manager - subservicesWatcher *services.FailureWatcher - hasSubservices bool - - backend resource.StorageBackend - cfg *setting.Cfg - features featuremgmt.FeatureToggles - db infraDB.DB - stopCh chan struct{} - stoppedCh chan error - - handler grpcserver.Provider - - tracing trace.Tracer - - authenticator func(ctx context.Context) (context.Context, error) - + backend resource.StorageBackend + cfg *setting.Cfg + features featuremgmt.FeatureToggles + stopCh chan struct{} + stoppedCh chan error + authenticator func(context.Context) (context.Context, error) + tracing trace.Tracer + db infraDB.DB log log.Logger reg prometheus.Registerer + docBuilders resource.DocumentBuilderSupplier storageMetrics *resource.StorageMetrics indexMetrics *resource.BleveIndexMetrics - - docBuilders resource.DocumentBuilderSupplier - searchRing *ring.Ring + + // Handler for the gRPC server + handler grpcserver.Provider + + // Ring lifecycle and sharding support ringLifecycler *ring.BasicLifecycler - queue QOSEnqueueDequeuer + // QoS support + queue *scheduler.Queue scheduler *scheduler.Scheduler + + // Subservices management + hasSubservices bool + subservices *services.Manager + subservicesWatcher *services.FailureWatcher } +// ProvideUnifiedStorageGrpcService provides a combined storage and search service running on the same gRPC server. +// This is used when running Grafana as a monolith where both services share the same process. +// Each service (storage and search) maintains its own lifecycle but shares the gRPC server. func ProvideUnifiedStorageGrpcService( cfg *setting.Cfg, features featuremgmt.FeatureToggles, @@ -100,7 +104,7 @@ func ProvideUnifiedStorageGrpcService( backend resource.StorageBackend, ) (UnifiedStorageGrpcService, error) { var err error - tracer := otel.Tracer("unified-storage") + tracer := otel.Tracer("unified-storage-combined") // FIXME: This is a temporary solution while we are migrating to the new authn interceptor // grpcutils.NewGrpcAuthenticator should be used instead. @@ -177,7 +181,7 @@ func ProvideUnifiedStorageGrpcService( MaxSizePerTenant: cfg.QOSMaxSizePerTenant, Registerer: qosReg, }) - scheduler, err := scheduler.NewScheduler(queue, &scheduler.Config{ + sched, err := scheduler.NewScheduler(queue, &scheduler.Config{ NumWorkers: cfg.QOSNumberWorker, Logger: log, }) @@ -186,7 +190,7 @@ func ProvideUnifiedStorageGrpcService( } s.queue = queue - s.scheduler = scheduler + s.scheduler = sched subservices = append(subservices, s.queue, s.scheduler) } @@ -199,6 +203,7 @@ func ProvideUnifiedStorageGrpcService( } // This will be used when running as a dskit service + // Note: We use StorageServer as the module name for backward compatibility s.BasicService = services.NewBasicService(s.starting, s.running, s.stopping).WithName(modules.StorageServer) return s, nil @@ -219,11 +224,6 @@ func (s *service) PrepareDownscale(w http.ResponseWriter, r *http.Request) { } } -var ( - // operation used by the search-servers to check if they own the namespace - searchOwnerRead = ring.NewOp([]ring.InstanceState{ring.JOINING, ring.ACTIVE, ring.LEAVING}, nil) -) - func (s *service) OwnsIndex(key resource.NamespacedResource) (bool, error) { if s.searchRing == nil { return true, nil @@ -260,71 +260,86 @@ func (s *service) starting(ctx context.Context) error { return err } - serverOptions := ServerOptions{ - Backend: s.backend, - DB: s.db, - Cfg: s.cfg, - Tracer: s.tracing, - Reg: s.reg, - AccessClient: authzClient, - StorageMetrics: s.storageMetrics, - Features: s.features, - QOSQueue: s.queue, - } - + // Setup overrides service if enabled + var overridesSvc *resource.OverridesService if s.cfg.OverridesFilePath != "" { - overridesSvc, err := resource.NewOverridesService(context.Background(), s.log, s.reg, s.tracing, resource.ReloadOptions{ + overridesSvc, err = resource.NewOverridesService(context.Background(), s.log, s.reg, s.tracing, resource.ReloadOptions{ FilePath: s.cfg.OverridesFilePath, ReloadPeriod: s.cfg.OverridesReloadInterval, }) if err != nil { return err } - serverOptions.OverridesService = overridesSvc } - // Handle search mode: "", "embedded", or "remote" - // Empty string defaults to "embedded" for backward compatibility - var searchServer resource.SearchServer - registerSearchServices := true - - switch s.cfg.SearchMode { - case "remote": - // Use remote search client - don't register search services locally - s.log.Info("Using remote search server", "address", s.cfg.SearchServerAddress) - remoteSearch, err := newRemoteSearchClient(s.cfg.SearchServerAddress) + // Ensure we have a backend - create one if needed + // This is critical: we create the backend ONCE and share it between search and storage servers + // to avoid duplicate metrics registration + backend := s.backend + if backend == nil { + eDB, err := dbimpl.ProvideResourceDB(s.db, s.cfg, s.tracing) if err != nil { - return fmt.Errorf("failed to create remote search client: %w", err) + return fmt.Errorf("failed to create resource DB: %w", err) } - searchServer = remoteSearch - registerSearchServices = false - case "", "embedded": - // Default: create local search server (backward compatible) - s.log.Info("Using embedded search server") - // SearchOptions are already configured, NewResourceServer will create the search server + isHA := isHighAvailabilityEnabled(s.cfg.SectionWithEnvOverrides("database"), + s.cfg.SectionWithEnvOverrides("resource_api")) - default: - return fmt.Errorf("invalid search_mode: %s (valid values: \"\", \"embedded\", \"remote\")", s.cfg.SearchMode) + b, err := NewBackend(BackendOptions{ + DBProvider: eDB, + Reg: s.reg, + IsHA: isHA, + storageMetrics: s.storageMetrics, + LastImportTimeMaxAge: s.cfg.MaxFileIndexAge, + }) + if err != nil { + return fmt.Errorf("failed to create backend: %w", err) + } + + // Initialize the backend + if err := b.Init(context.Background()); err != nil { + return fmt.Errorf("failed to initialize backend: %w", err) + } + backend = b } - var server resource.ResourceServer - if searchServer != nil { - // Remote search mode: pass the remote search client to the resource server - // Clear local search options since we're using remote search - serverOptions.SearchOptions = resource.SearchOptions{} - var err error - server, _, err = NewResourceServer(serverOptions) - if err != nil { - return err - } - } else { - // Embedded mode: create both server and search server together - var err error - server, searchServer, err = NewResourceServer(serverOptions) - if err != nil { - return err - } + // Create search options for the search server + searchOptions, err := search.NewSearchOptions(s.features, s.cfg, s.docBuilders, s.indexMetrics, s.OwnsIndex) + if err != nil { + return err + } + + // Create the search server - pass the shared backend + searchServer, err := NewSearchServer(SearchServerOptions{ + Backend: backend, // Use the shared backend + DB: s.db, + Cfg: s.cfg, + Tracer: s.tracing, + Reg: s.reg, + AccessClient: authzClient, + SearchOptions: searchOptions, + IndexMetrics: s.indexMetrics, + OwnsIndexFn: s.OwnsIndex, + }) + if err != nil { + return err + } + + // Create the storage server - pass the shared backend + storageServer, err := NewStorageServer(StorageServerOptions{ + Backend: backend, // Use the shared backend + OverridesService: overridesSvc, + DB: s.db, + Cfg: s.cfg, + Tracer: s.tracing, + Reg: s.reg, + AccessClient: authzClient, + StorageMetrics: s.storageMetrics, + Features: s.features, + QOSQueue: s.queue, + }) + if err != nil { + return err } s.handler, err = grpcserver.ProvideService(s.cfg, s.features, interceptors.AuthenticatorFunc(s.authenticator), s.tracing, prometheus.DefaultRegisterer) @@ -332,22 +347,21 @@ func (s *service) starting(ctx context.Context) error { return err } - healthService, err := resource.ProvideHealthService(server) + healthService, err := resource.ProvideHealthService(storageServer) if err != nil { return err } srv := s.handler.GetServer() - resourcepb.RegisterResourceStoreServer(srv, server) - resourcepb.RegisterBulkStoreServer(srv, server) - // Only register search services if running in embedded mode - if registerSearchServices { - resourcepb.RegisterResourceIndexServer(srv, searchServer) - resourcepb.RegisterManagedObjectIndexServer(srv, searchServer) - } - resourcepb.RegisterBlobStoreServer(srv, server) - resourcepb.RegisterDiagnosticsServer(srv, server) - resourcepb.RegisterQuotasServer(srv, server) + // Register storage services + resourcepb.RegisterResourceStoreServer(srv, storageServer) + resourcepb.RegisterBulkStoreServer(srv, storageServer) + resourcepb.RegisterBlobStoreServer(srv, storageServer) + resourcepb.RegisterDiagnosticsServer(srv, storageServer) + resourcepb.RegisterQuotasServer(srv, storageServer) + // Register search services + resourcepb.RegisterResourceIndexServer(srv, searchServer) + resourcepb.RegisterManagedObjectIndexServer(srv, searchServer) grpc_health_v1.RegisterHealthServer(srv, healthService) // register reflection service diff --git a/pkg/storage/unified/sql/storage_service.go b/pkg/storage/unified/sql/storage_service.go new file mode 100644 index 00000000000..c7d04c5ef43 --- /dev/null +++ b/pkg/storage/unified/sql/storage_service.go @@ -0,0 +1,234 @@ +package sql + +import ( + "context" + "errors" + "fmt" + + "github.com/prometheus/client_golang/prometheus" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/trace" + "google.golang.org/grpc/health/grpc_health_v1" + + "github.com/grafana/dskit/services" + + infraDB "github.com/grafana/grafana/pkg/infra/db" + "github.com/grafana/grafana/pkg/infra/log" + "github.com/grafana/grafana/pkg/modules" + "github.com/grafana/grafana/pkg/services/authz" + "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/resource" + "github.com/grafana/grafana/pkg/storage/unified/resource/grpc" + "github.com/grafana/grafana/pkg/storage/unified/resourcepb" + "github.com/grafana/grafana/pkg/util/scheduler" +) + +var _ UnifiedStorageGrpcService = (*storageService)(nil) + +type storageService struct { + *services.BasicService + + backend resource.StorageBackend + cfg *setting.Cfg + features featuremgmt.FeatureToggles + db infraDB.DB + stopCh chan struct{} + stoppedCh chan error + + handler grpcserver.Provider + + tracing trace.Tracer + + authenticator func(ctx context.Context) (context.Context, error) + + log log.Logger + reg prometheus.Registerer + storageMetrics *resource.StorageMetrics + + queue QOSEnqueueDequeuer + scheduler *scheduler.Scheduler + + // Subservices manager + subservices *services.Manager + subservicesWatcher *services.FailureWatcher + hasSubservices bool +} + +func ProvideStorageService( + cfg *setting.Cfg, + features featuremgmt.FeatureToggles, + db infraDB.DB, + log log.Logger, + reg prometheus.Registerer, + storageMetrics *resource.StorageMetrics, + backend resource.StorageBackend, +) (UnifiedStorageGrpcService, error) { + var err error + tracer := otel.Tracer("unified-storage-server") + + authn := NewAuthenticatorWithFallback(cfg, reg, tracer, func(ctx context.Context) (context.Context, error) { + auth := grpc.Authenticator{Tracer: tracer} + return auth.Authenticate(ctx) + }) + + s := &storageService{ + backend: backend, + cfg: cfg, + features: features, + stopCh: make(chan struct{}), + stoppedCh: make(chan error, 1), + authenticator: authn, + tracing: tracer, + db: db, + log: log, + reg: reg, + storageMetrics: storageMetrics, + subservicesWatcher: services.NewFailureWatcher(), + } + + subservices := []services.Service{} + + // Setup QOS if enabled + if cfg.QOSEnabled { + qosReg := prometheus.WrapRegistererWithPrefix("resource_server_qos_", reg) + queue := scheduler.NewQueue(&scheduler.QueueOptions{ + MaxSizePerTenant: cfg.QOSMaxSizePerTenant, + Registerer: qosReg, + }) + sched, err := scheduler.NewScheduler(queue, &scheduler.Config{ + NumWorkers: cfg.QOSNumberWorker, + Logger: log, + }) + if err != nil { + return nil, fmt.Errorf("failed to create qos scheduler: %s", err) + } + + s.queue = queue + s.scheduler = sched + subservices = append(subservices, s.queue, s.scheduler) + } + + if len(subservices) > 0 { + s.hasSubservices = true + s.subservices, err = services.NewManager(subservices...) + if err != nil { + return nil, fmt.Errorf("failed to create subservices manager: %w", err) + } + } + + // This will be used when running as a dskit service + s.BasicService = services.NewBasicService(s.starting, s.running, s.stopping).WithName(modules.StorageServer) + + return s, nil +} + +func (s *storageService) starting(ctx context.Context) error { + if s.hasSubservices { + s.subservicesWatcher.WatchManager(s.subservices) + if err := services.StartManagerAndAwaitHealthy(ctx, s.subservices); err != nil { + return fmt.Errorf("failed to start subservices: %w", err) + } + } + + authzClient, err := authz.ProvideStandaloneAuthZClient(s.cfg, s.features, s.tracing, s.reg) + if err != nil { + return err + } + + // Setup overrides service if enabled + var overridesSvc *resource.OverridesService + if s.cfg.OverridesFilePath != "" { + overridesSvc, err = resource.NewOverridesService(context.Background(), s.log, s.reg, s.tracing, resource.ReloadOptions{ + FilePath: s.cfg.OverridesFilePath, + ReloadPeriod: s.cfg.OverridesReloadInterval, + }) + if err != nil { + return err + } + } + + // Create the storage server + storageServer, err := NewStorageServer(StorageServerOptions{ + Backend: s.backend, + OverridesService: overridesSvc, + DB: s.db, + Cfg: s.cfg, + Tracer: s.tracing, + Reg: s.reg, + AccessClient: authzClient, + StorageMetrics: s.storageMetrics, + Features: s.features, + QOSQueue: s.queue, + }) + if err != nil { + return err + } + + s.handler, err = grpcserver.ProvideService(s.cfg, s.features, interceptors.AuthenticatorFunc(s.authenticator), s.tracing, prometheus.DefaultRegisterer) + if err != nil { + return err + } + + healthService, err := resource.ProvideHealthService(storageServer) + if err != nil { + return err + } + + srv := s.handler.GetServer() + // Register storage services + resourcepb.RegisterResourceStoreServer(srv, storageServer) + resourcepb.RegisterBulkStoreServer(srv, storageServer) + resourcepb.RegisterBlobStoreServer(srv, storageServer) + resourcepb.RegisterDiagnosticsServer(srv, storageServer) + resourcepb.RegisterQuotasServer(srv, storageServer) + grpc_health_v1.RegisterHealthServer(srv, healthService) + + // register reflection service + _, err = grpcserver.ProvideReflectionService(s.cfg, s.handler) + if err != nil { + return err + } + + // start the gRPC server + go func() { + err := s.handler.Run(ctx) + if err != nil { + s.stoppedCh <- err + } else { + s.stoppedCh <- nil + } + }() + return nil +} + +// GetAddress returns the address of the gRPC server. +func (s *storageService) GetAddress() string { + return s.handler.GetAddress() +} + +func (s *storageService) running(ctx context.Context) error { + select { + case err := <-s.stoppedCh: + if err != nil && !errors.Is(err, context.Canceled) { + return err + } + case err := <-s.subservicesWatcher.Chan(): + return fmt.Errorf("subservice failure: %w", err) + case <-ctx.Done(): + close(s.stopCh) + } + return nil +} + +func (s *storageService) stopping(_ error) error { + if s.hasSubservices { + err := services.StopManagerAndAwaitStopped(context.Background(), s.subservices) + if err != nil { + return fmt.Errorf("failed to stop subservices: %w", err) + } + } + return nil +} diff --git a/pkg/storage/unified/testing/search_and_storage.go b/pkg/storage/unified/testing/search_and_storage.go index 0f12b422031..9b2ec03ee57 100644 --- a/pkg/storage/unified/testing/search_and_storage.go +++ b/pkg/storage/unified/testing/search_and_storage.go @@ -115,10 +115,9 @@ func RunTestSearchAndStorage(t *testing.T, ctx context.Context, backend resource err = searchServer.Init(ctx) require.NoError(t, err) - // Create a resource server with the search server + // Create a resource server separately from the search server server, err = resource.NewResourceServer(resource.ResourceServerOptions{ Backend: backend, - Search: searchServer, }) require.NoError(t, err) })