From a7aa55f9082a564e2bd9d1dbad3eeaa7043d1443 Mon Sep 17 00:00:00 2001 From: mayor Date: Wed, 14 Jan 2026 17:49:03 +0100 Subject: [PATCH] Add search-server target and configurable search mode - Add SearchServer module target for standalone search service - Add search_mode config: "", "embedded" (default), "remote" - Add search_server_address config for remote search server - Create remote_search.go: gRPC client wrapper for remote search - Create search_service.go: standalone search gRPC service - Modify service.go: conditional search mode handling - Backward compatible: empty/embedded mode works as before Co-Authored-By: Claude Opus 4.5 --- pkg/modules/dependencies.go | 2 + pkg/server/module_server.go | 8 + pkg/setting/setting.go | 2 + pkg/setting/setting_unified_storage.go | 4 + pkg/storage/unified/sql/remote_search.go | 83 +++++ pkg/storage/unified/sql/search_service.go | 353 ++++++++++++++++++++++ pkg/storage/unified/sql/service.go | 53 +++- 7 files changed, 500 insertions(+), 5 deletions(-) create mode 100644 pkg/storage/unified/sql/remote_search.go create mode 100644 pkg/storage/unified/sql/search_service.go diff --git a/pkg/modules/dependencies.go b/pkg/modules/dependencies.go index d78240a3434..07d7888f8f5 100644 --- a/pkg/modules/dependencies.go +++ b/pkg/modules/dependencies.go @@ -10,6 +10,7 @@ const ( SearchServerRing string = "search-server-ring" SearchServerDistributor string = "search-server-distributor" StorageServer string = "storage-server" + SearchServer string = "search-server" ZanzanaServer string = "zanzana-server" InstrumentationServer string = "instrumentation-server" FrontendServer string = "frontend-server" @@ -21,6 +22,7 @@ var dependencyMap = map[string][]string{ SearchServerRing: {InstrumentationServer, MemberlistKV}, GrafanaAPIServer: {InstrumentationServer}, StorageServer: {InstrumentationServer, SearchServerRing}, + SearchServer: {InstrumentationServer, SearchServerRing}, ZanzanaServer: {InstrumentationServer}, SearchServerDistributor: {InstrumentationServer, MemberlistKV, SearchServerRing}, Core: {}, diff --git a/pkg/server/module_server.go b/pkg/server/module_server.go index 3954b29dac9..49d655c0e74 100644 --- a/pkg/server/module_server.go +++ b/pkg/server/module_server.go @@ -205,6 +205,14 @@ func (s *ModuleServer) Run() error { return sql.ProvideUnifiedStorageGrpcService(s.cfg, s.features, nil, s.log, s.registerer, docBuilders, s.storageMetrics, s.indexMetrics, s.searchServerRing, s.MemberlistKVConfig, s.httpServerRouter, s.storageBackend) }) + m.RegisterModule(modules.SearchServer, func() (services.Service, error) { + docBuilders, err := InitializeDocumentBuilders(s.cfg) + if err != nil { + return nil, err + } + return sql.ProvideSearchGrpcService(s.cfg, s.features, nil, s.log, s.registerer, docBuilders, s.indexMetrics, s.searchServerRing, s.MemberlistKVConfig, s.storageBackend) + }) + m.RegisterModule(modules.ZanzanaServer, func() (services.Service, error) { return authz.ProvideZanzanaService(s.cfg, s.features, s.registerer) }) diff --git a/pkg/setting/setting.go b/pkg/setting/setting.go index 9667b82b9fa..b16f3a3e2ca 100644 --- a/pkg/setting/setting.go +++ b/pkg/setting/setting.go @@ -623,6 +623,8 @@ 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 21a3f455993..fe1f021a6c0 100644 --- a/pkg/setting/setting_unified_storage.go +++ b/pkg/setting/setting_unified_storage.go @@ -136,6 +136,10 @@ 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/sql/remote_search.go b/pkg/storage/unified/sql/remote_search.go new file mode 100644 index 00000000000..beff31d70d5 --- /dev/null +++ b/pkg/storage/unified/sql/remote_search.go @@ -0,0 +1,83 @@ +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 +} + +// 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), + }, 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) +} diff --git a/pkg/storage/unified/sql/search_service.go b/pkg/storage/unified/sql/search_service.go new file mode 100644 index 00000000000..ea1874d0c86 --- /dev/null +++ b/pkg/storage/unified/sql/search_service.go @@ -0,0 +1,353 @@ +package sql + +import ( + "context" + "errors" + "fmt" + "hash/fnv" + "net" + "os" + "strconv" + "time" + + "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" + + 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/storage/unified/search" +) + +var ( + _ SearchGrpcService = (*searchService)(nil) +) + +// SearchGrpcService is the interface for the standalone search gRPC service. +type SearchGrpcService interface { + services.NamedService + + // GetAddress returns the address where this service is running + GetAddress() string +} + +type searchService struct { + *services.BasicService + + // 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 +} + +// ProvideSearchGrpcService creates a standalone search gRPC service. +// This is used when running search-server as a separate target. +func ProvideSearchGrpcService( + cfg *setting.Cfg, + features featuremgmt.FeatureToggles, + db infraDB.DB, + log log.Logger, + reg prometheus.Registerer, + docBuilders resource.DocumentBuilderSupplier, + indexMetrics *resource.BleveIndexMetrics, + searchRing *ring.Ring, + memberlistKVConfig kv.Config, + backend resource.StorageBackend, +) (SearchGrpcService, error) { + tracer := otel.Tracer("search-server") + + authn := NewAuthenticatorWithFallback(cfg, reg, tracer, func(ctx context.Context) (context.Context, error) { + auth := grpc.Authenticator{Tracer: tracer} + return auth.Authenticate(ctx) + }) + + s := &searchService{ + cfg: cfg, + features: features, + stopCh: make(chan struct{}), + stoppedCh: make(chan error, 1), + authenticator: authn, + tracing: tracer, + db: db, + log: log, + reg: reg, + docBuilders: docBuilders, + indexMetrics: indexMetrics, + searchRing: searchRing, + subservicesWatcher: services.NewFailureWatcher(), + backend: backend, + } + + subservices := []services.Service{} + if cfg.EnableSharding { + ringStore, err := kv.NewClient( + memberlistKVConfig, + ring.GetCodec(), + kv.RegistererWithKVName(reg, resource.RingName), + log, + ) + if err != nil { + return nil, fmt.Errorf("failed to create KV store client: %s", err) + } + + lifecyclerCfg, err := toSearchLifecyclerConfig(cfg, log) + if err != nil { + return nil, fmt.Errorf("failed to initialize search-ring lifecycler config: %s", err) + } + + delegate := ring.BasicLifecyclerDelegate(ring.NewInstanceRegisterDelegate(ring.JOINING, resource.RingNumTokens)) + delegate = ring.NewLeaveOnStoppingDelegate(delegate, log) + delegate = ring.NewAutoForgetDelegate(resource.RingHeartbeatTimeout*2, delegate, log) + + s.ringLifecycler, err = ring.NewBasicLifecycler( + lifecyclerCfg, + resource.RingName, + resource.RingKey, + ringStore, + delegate, + log, + reg, + ) + if err != nil { + return nil, fmt.Errorf("failed to initialize search-ring lifecycler: %s", err) + } + + s.ringLifecycler.SetKeepInstanceInTheRingOnShutdown(true) + subservices = append(subservices, s.ringLifecycler) + } + + 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) + } + } + + s.BasicService = services.NewBasicService(s.starting, s.running, s.stopping).WithName(modules.SearchServer) + + return s, nil +} + +func (s *searchService) OwnsIndex(key resource.NamespacedResource) (bool, error) { + if s.searchRing == nil { + return true, nil + } + + if st := s.searchRing.State(); st != services.Running { + return false, fmt.Errorf("ring is not Running: %s", st) + } + + ringHasher := fnv.New32a() + _, err := ringHasher.Write([]byte(key.Namespace)) + if err != nil { + return false, fmt.Errorf("error hashing namespace: %w", err) + } + + rs, err := s.searchRing.GetWithOptions(ringHasher.Sum32(), searchOwnerRead, ring.WithReplicationFactor(s.searchRing.ReplicationFactor())) + if err != nil { + return false, fmt.Errorf("error getting replicaset from ring: %w", err) + } + + return rs.Includes(s.ringLifecycler.GetInstanceAddr()), nil +} + +func (s *searchService) 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 + } + + 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) + 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) + } + + s.handler, err = grpcserver.ProvideService(s.cfg, s.features, interceptors.AuthenticatorFunc(s.authenticator), s.tracing, prometheus.DefaultRegisterer) + if err != nil { + return err + } + + srv := s.handler.GetServer() + resourcepb.RegisterResourceIndexServer(srv, searchServer) + resourcepb.RegisterManagedObjectIndexServer(srv, searchServer) + grpc_health_v1.RegisterHealthServer(srv, &searchHealthService{searchServer: searchServer}) + + // register reflection service + _, err = grpcserver.ProvideReflectionService(s.cfg, s.handler) + if err != nil { + return err + } + + if s.cfg.EnableSharding { + s.log.Info("waiting until search server is JOINING in the ring") + lfcCtx, cancel := context.WithTimeout(context.Background(), s.cfg.ResourceServerJoinRingTimeout) + defer cancel() + if err := ring.WaitInstanceState(lfcCtx, s.searchRing, s.ringLifecycler.GetInstanceID(), ring.JOINING); err != nil { + return fmt.Errorf("error switching to JOINING in the ring: %s", err) + } + s.log.Info("search server is JOINING in the ring") + + if err := s.ringLifecycler.ChangeState(ctx, ring.ACTIVE); err != nil { + return fmt.Errorf("error switching to ACTIVE in the ring: %s", err) + } + s.log.Info("search server is ACTIVE in the ring") + } + + // start the gRPC server + go func() { + err := s.handler.Run(ctx) + if err != nil { + s.stoppedCh <- err + } else { + s.stoppedCh <- nil + } + }() + return nil +} + +func (s *searchService) GetAddress() string { + return s.handler.GetAddress() +} + +func (s *searchService) 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 *searchService) 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 +} + +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 + } + + 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 +} + +// searchHealthService implements the health check for the search service. +type searchHealthService struct { + searchServer resource.SearchServer +} + +func (h *searchHealthService) Check(ctx context.Context, req *grpc_health_v1.HealthCheckRequest) (*grpc_health_v1.HealthCheckResponse, error) { + return &grpc_health_v1.HealthCheckResponse{ + Status: grpc_health_v1.HealthCheckResponse_SERVING, + }, nil +} + +func (h *searchHealthService) Watch(req *grpc_health_v1.HealthCheckRequest, server grpc_health_v1.Health_WatchServer) error { + return fmt.Errorf("watch not implemented") +} + +func (h *searchHealthService) List(ctx context.Context, req *grpc_health_v1.HealthListRequest) (*grpc_health_v1.HealthListResponse, error) { + check, err := h.Check(ctx, &grpc_health_v1.HealthCheckRequest{}) + if err != nil { + return nil, err + } + return &grpc_health_v1.HealthListResponse{ + Statuses: map[string]*grpc_health_v1.HealthCheckResponse{ + "": check, + }, + }, nil +} diff --git a/pkg/storage/unified/sql/service.go b/pkg/storage/unified/sql/service.go index 9bee17f445d..bbf8c4580ef 100644 --- a/pkg/storage/unified/sql/service.go +++ b/pkg/storage/unified/sql/service.go @@ -292,10 +292,50 @@ func (s *service) starting(ctx context.Context) error { serverOptions.OverridesService = overridesSvc } - server, searchServer, err := NewResourceServer(serverOptions) - if err != nil { - return err + // 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) + if err != nil { + return fmt.Errorf("failed to create remote search client: %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 + + default: + return fmt.Errorf("invalid search_mode: %s (valid values: \"\", \"embedded\", \"remote\")", s.cfg.SearchMode) } + + 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 + } + } + s.handler, err = grpcserver.ProvideService(s.cfg, s.features, interceptors.AuthenticatorFunc(s.authenticator), s.tracing, prometheus.DefaultRegisterer) if err != nil { return err @@ -309,8 +349,11 @@ func (s *service) starting(ctx context.Context) error { srv := s.handler.GetServer() resourcepb.RegisterResourceStoreServer(srv, server) resourcepb.RegisterBulkStoreServer(srv, server) - resourcepb.RegisterResourceIndexServer(srv, searchServer) - resourcepb.RegisterManagedObjectIndexServer(srv, searchServer) + // 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)