package options import ( "context" "fmt" "net" "strconv" "strings" "time" "github.com/spf13/pflag" "go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc" "google.golang.org/grpc" "google.golang.org/grpc/backoff" "google.golang.org/grpc/codes" "google.golang.org/grpc/credentials/insecure" "google.golang.org/grpc/keepalive" genericapiserver "k8s.io/apiserver/pkg/server" "k8s.io/apiserver/pkg/server/options" "k8s.io/client-go/rest" grpc_retry "github.com/grpc-ecosystem/go-grpc-middleware/retry" apiserverrest "github.com/grafana/grafana/pkg/apiserver/rest" "github.com/grafana/grafana/pkg/infra/tracing" secret "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/authn/grpcutils" "github.com/grafana/grafana/pkg/setting" "github.com/grafana/grafana/pkg/storage/unified/apistore" "github.com/grafana/grafana/pkg/storage/unified/resource" ) type StorageType string const ( StorageTypeFile StorageType = "file" StorageTypeEtcd StorageType = "etcd" StorageTypeUnified StorageType = "unified" StorageTypeUnifiedGrpc StorageType = "unified-grpc" StorageTypeUnifiedKVGrpc StorageType = "unified-kv-grpc" // Deprecated: legacy is a shim that is no longer necessary StorageTypeLegacy StorageType = "legacy" BlobThresholdDefault int = 0 ) type RestConfigProvider interface { GetRestConfig(context.Context) (*rest.Config, error) } type StorageOptions struct { // The desired storage type StorageType StorageType // For unified-grpc Address string SearchServerAddress string GrpcClientAuthenticationToken string GrpcClientAuthenticationTokenExchangeURL string GrpcClientAuthenticationTokenNamespace string GrpcClientAuthenticationAllowInsecure bool GrpcClientKeepaliveTime time.Duration // Secrets Manager Configuration for InlineSecureValueSupport SecretsManagerGrpcClientEnable bool SecretsManagerGrpcClientLoadBalancing bool SecretsManagerGrpcServerAddress string SecretsManagerGrpcServerUseTLS bool SecretsManagerGrpcServerTLSSkipVerify bool SecretsManagerGrpcServerTLSServerName string SecretsManagerGrpcServerTLSCAFile string // For file storage, this is the requested path DataPath string // Optional blob storage connection string // file:///path/to/dir // gs://my-bucket (using default credentials) // s3://my-bucket?region=us-west-1 (using default credentials) // azblob://my-container BlobStoreURL string // Optional blob storage field. When an object's size in bytes exceeds the threshold // value, it is considered large and gets partially stored in blob storage. BlobThresholdBytes int // Support writing secrets inline InlineSecrets secret.InlineSecureValueSupport // {resource}.{group} = 1|2|3|4 UnifiedStorageConfig map[string]setting.UnifiedStorageConfig // Access to the other clients ConfigProvider RestConfigProvider } // unifiedStorageConfigValue implements pflag.Value for parsing unified storage config type unifiedStorageConfigValue struct { config *map[string]setting.UnifiedStorageConfig } func (v *unifiedStorageConfigValue) String() string { if v.config == nil || len(*v.config) == 0 { return "" } parts := make([]string, 0, len(*v.config)) for key, cfg := range *v.config { parts = append(parts, fmt.Sprintf("%s=%d", key, cfg.DualWriterMode)) } return strings.Join(parts, ",") } func (v *unifiedStorageConfigValue) Set(val string) error { if val == "" { return nil } // Parse comma-separated key=value pairs pairs := strings.Split(val, ",") for _, pair := range pairs { kv := strings.SplitN(pair, "=", 2) if len(kv) != 2 { return fmt.Errorf("invalid format: %s (expected key=value)", pair) } key := strings.TrimSpace(kv[0]) mode, err := strconv.Atoi(strings.TrimSpace(kv[1])) if err != nil { return fmt.Errorf("invalid mode value for %s: %w", key, err) } if mode < 0 || mode > 5 { return fmt.Errorf("mode must be between 0 and 5, got %d for %s", mode, key) } (*v.config)[key] = setting.UnifiedStorageConfig{ DualWriterMode: apiserverrest.DualWriterMode(mode), DualWriterMigrationDataSyncDisabled: true, } } return nil } func (v *unifiedStorageConfigValue) Type() string { return "stringToUnifiedStorageConfig" } func NewStorageOptions() *StorageOptions { return &StorageOptions{ StorageType: StorageTypeUnified, Address: "localhost:10000", GrpcClientAuthenticationTokenNamespace: "*", GrpcClientAuthenticationAllowInsecure: false, GrpcClientKeepaliveTime: 0, BlobThresholdBytes: BlobThresholdDefault, UnifiedStorageConfig: make(map[string]setting.UnifiedStorageConfig), } } func (o *StorageOptions) AddFlags(fs *pflag.FlagSet) { fs.StringVar((*string)(&o.StorageType), "grafana-apiserver-storage-type", string(o.StorageType), "Storage type") fs.StringVar(&o.DataPath, "grafana-apiserver-storage-path", o.DataPath, "Storage path for file storage") fs.StringVar(&o.Address, "grafana-apiserver-storage-address", o.Address, "Remote grpc address endpoint") fs.StringVar(&o.SearchServerAddress, "grafana-apiserver-search-address", o.SearchServerAddress, "Remote grpc address endpoint for search server") fs.StringVar(&o.GrpcClientAuthenticationToken, "grpc-client-authentication-token", o.GrpcClientAuthenticationToken, "Token for grpc client authentication") fs.StringVar(&o.GrpcClientAuthenticationTokenExchangeURL, "grpc-client-authentication-token-exchange-url", o.GrpcClientAuthenticationTokenExchangeURL, "Token exchange url for grpc client authentication") fs.StringVar(&o.GrpcClientAuthenticationTokenNamespace, "grpc-client-authentication-token-namespace", o.GrpcClientAuthenticationTokenNamespace, "Token namespace for grpc client authentication") fs.BoolVar(&o.GrpcClientAuthenticationAllowInsecure, "grpc-client-authentication-allow-insecure", o.GrpcClientAuthenticationAllowInsecure, "Allow insecure grpc client authentication") fs.DurationVar(&o.GrpcClientKeepaliveTime, "grpc-client-keepalive-time", o.GrpcClientKeepaliveTime, "gRPC client keep-alive ping interval (e.g., 6m).") // Use custom flag value for unified storage config fs.Var(&unifiedStorageConfigValue{config: &o.UnifiedStorageConfig}, "grafana-apiserver-unified-storage-config", "Unified storage configuration per resource.group in the format resource.group=mode,... where mode is 0-5") // Secrets Manager Configuration flags fs.BoolVar(&o.SecretsManagerGrpcClientEnable, "grafana.secrets-manager.grpc-client-enable", false, "Enable gRPC client for secrets manager") fs.StringVar(&o.SecretsManagerGrpcServerAddress, "grafana.secrets-manager.grpc-server-address", "", "gRPC server address for secrets manager") fs.BoolVar(&o.SecretsManagerGrpcServerUseTLS, "grafana.secrets-manager.grpc-server-use-tls", false, "Use TLS for gRPC server communication") fs.BoolVar(&o.SecretsManagerGrpcServerTLSSkipVerify, "grafana.secrets-manager.grpc-server-tls-skip-verify", false, "Skip TLS verification for gRPC server") fs.StringVar(&o.SecretsManagerGrpcServerTLSServerName, "grafana.secrets-manager.grpc-server-tls-server-name", "", "Server name for TLS verification") fs.StringVar(&o.SecretsManagerGrpcServerTLSCAFile, "grafana.secrets-manager.grpc-server-tls-ca-file", "", "CA file for TLS verification") fs.BoolVar(&o.SecretsManagerGrpcClientLoadBalancing, "grafana.secrets-manager.grpc-client-load-balancing", false, "Enable client-side load balancing for gRPC client") } func (o *StorageOptions) Validate() []error { errs := []error{} switch o.StorageType { // nolint:staticcheck case StorageTypeLegacy: // no-op case StorageTypeUnifiedKVGrpc: // no-op (enterprise only) case StorageTypeFile, StorageTypeEtcd, StorageTypeUnified, StorageTypeUnifiedGrpc: // no-op default: // nolint:staticcheck errs = append(errs, fmt.Errorf("--grafana-apiserver-storage-type must be one of %s, %s, %s, %s, %s", StorageTypeFile, StorageTypeEtcd, StorageTypeLegacy, StorageTypeUnified, StorageTypeUnifiedGrpc)) } if _, _, err := net.SplitHostPort(o.Address); err != nil { errs = append(errs, fmt.Errorf("--grafana-apiserver-storage-address must be a valid network address: %v", err)) } // Only works for single tenant grafana right now if o.BlobStoreURL != "" && o.StorageType != StorageTypeUnified { errs = append(errs, fmt.Errorf("blob storage is only valid with unified storage")) } // Validate grpc client with auth if o.StorageType == StorageTypeUnifiedGrpc && o.GrpcClientAuthenticationToken != "" { if o.GrpcClientAuthenticationToken == "" { errs = append(errs, fmt.Errorf("grpc client auth token is required for unified-grpc storage")) } if o.GrpcClientAuthenticationTokenExchangeURL == "" { errs = append(errs, fmt.Errorf("grpc client auth token exchange url is required for unified-grpc storage")) } if o.GrpcClientAuthenticationTokenNamespace == "" { errs = append(errs, fmt.Errorf("grpc client auth namespace is required for unified-grpc storage")) } } if o.SecretsManagerGrpcClientEnable { if o.SecretsManagerGrpcServerAddress == "" { errs = append(errs, fmt.Errorf("secrets manager grpc server address is required for secrets manager grpc client")) } if o.SecretsManagerGrpcServerUseTLS && !o.SecretsManagerGrpcServerTLSSkipVerify && o.SecretsManagerGrpcServerTLSCAFile == "" { errs = append(errs, fmt.Errorf("secrets manager grpc server ca file is required for secrets manager grpc client")) } } return errs } func (o *StorageOptions) ApplyTo(serverConfig *genericapiserver.RecommendedConfig, etcdOptions *options.EtcdOptions, tracer tracing.Tracer, secureServing *options.SecureServingOptions) error { if o.StorageType != StorageTypeUnifiedGrpc { return nil } grpcOpts := o.buildGrpcDialOptions() conn, err := grpc.NewClient(o.Address, grpcOpts...) if err != nil { return err } var indexConn *grpc.ClientConn if o.SearchServerAddress != "" { indexConn, err = grpc.NewClient(o.SearchServerAddress, grpcOpts...) if err != nil { return err } } else { indexConn = conn } const resourceStoreAudience = "resourceStore" unified, err := resource.NewRemoteResourceClient(tracer, conn, indexConn, resource.RemoteResourceClientConfig{ Token: o.GrpcClientAuthenticationToken, TokenExchangeURL: o.GrpcClientAuthenticationTokenExchangeURL, Namespace: o.GrpcClientAuthenticationTokenNamespace, Audiences: []string{resourceStoreAudience}, }) if err != nil { return err } // setup inline secrets if configured if o.InlineSecrets == nil && o.SecretsManagerGrpcClientEnable { tlsCfg := inlinesecurevalue.TLSConfig{ UseTLS: o.SecretsManagerGrpcServerUseTLS, CAFile: o.SecretsManagerGrpcServerTLSCAFile, ServerName: o.SecretsManagerGrpcServerTLSServerName, InsecureSkipVerify: o.SecretsManagerGrpcServerTLSSkipVerify, } inlineSecureValueService, err := inlinesecurevalue.NewGRPCSecureValueService( &grpcutils.GrpcClientConfig{ Token: o.GrpcClientAuthenticationToken, TokenExchangeURL: o.GrpcClientAuthenticationTokenExchangeURL, TokenNamespace: o.GrpcClientAuthenticationTokenNamespace, }, o.SecretsManagerGrpcServerAddress, tlsCfg, tracer, o.SecretsManagerGrpcClientLoadBalancing, ) if err != nil { return fmt.Errorf("failed to create inline secure value service: %w", err) } o.InlineSecrets = inlineSecureValueService } getter := apistore.NewRESTOptionsGetterForClient(unified, o.InlineSecrets, etcdOptions.StorageConfig, o.ConfigProvider) serverConfig.RESTOptionsGetter = getter return nil } // buildGrpcDialOptions creates gRPC dial options with resilience mechanisms: // - Round-robin load balancing with client-side health checking // - Retry interceptor for transient connection issues // - Keepalive for long-lived connections func (o *StorageOptions) buildGrpcDialOptions() []grpc.DialOption { // Retry interceptor for transient connection issues (codes.Unavailable includes connection refused) retryInterceptor := grpc_retry.UnaryClientInterceptor( grpc_retry.WithMax(3), grpc_retry.WithBackoff(grpc_retry.BackoffExponentialWithJitter(time.Second, 0.5)), grpc_retry.WithCodes(codes.ResourceExhausted, codes.Unavailable), ) opts := []grpc.DialOption{ grpc.WithStatsHandler(otelgrpc.NewClientHandler()), grpc.WithTransportCredentials(insecure.NewCredentials()), grpc.WithChainUnaryInterceptor(retryInterceptor), grpc.WithDefaultServiceConfig(`{"loadBalancingPolicy":"round_robin"}`), grpc.WithConnectParams(grpc.ConnectParams{ Backoff: backoff.Config{ BaseDelay: 100 * time.Millisecond, Multiplier: 1.6, Jitter: 0.2, MaxDelay: 10 * time.Second, }, MinConnectTimeout: 5 * time.Second, }), } if o.GrpcClientKeepaliveTime > 0 { opts = append(opts, grpc.WithKeepaliveParams(keepalive.ClientParameters{ Time: o.GrpcClientKeepaliveTime, Timeout: 10 * time.Second, PermitWithoutStream: true, })) } return opts }