Provisioning: route provisioning storage to the provisioning client (#104082)

This commit is contained in:
Ryan McKinley
2025-04-18 21:47:37 +03:00
committed by GitHub
parent 689da86f81
commit efd9334295
9 changed files with 151 additions and 49 deletions
@@ -40,7 +40,7 @@ func (s *DashboardStorage) NewStore(dash utils.ResourceInfo, scheme *runtime.Sch
}
client := legacy.NewDirectResourceClient(server) // same context
optsGetter := apistore.NewRESTOptionsGetterForClient(client,
defaultOpts.StorageConfig.Config,
defaultOpts.StorageConfig.Config, nil,
)
optsGetter.RegisterOptions(dash.GroupResource(), apistore.StorageOptions{
EnableFolderSupport: true,
+10 -1
View File
@@ -1,6 +1,7 @@
package options
import (
"context"
"fmt"
"net"
@@ -10,6 +11,7 @@ import (
"google.golang.org/grpc/credentials/insecure"
genericapiserver "k8s.io/apiserver/pkg/server"
"k8s.io/apiserver/pkg/server/options"
"k8s.io/client-go/rest"
"github.com/grafana/grafana/pkg/infra/tracing"
"github.com/grafana/grafana/pkg/services/featuremgmt"
@@ -32,6 +34,10 @@ const (
BlobThresholdDefault int = 0
)
type RestConfigProvider interface {
GetRestConfig(context.Context) (*rest.Config, error)
}
type StorageOptions struct {
// The desired storage type
StorageType StorageType
@@ -58,6 +64,9 @@ type StorageOptions struct {
// {resource}.{group} = 1|2|3|4
UnifiedStorageConfig map[string]setting.UnifiedStorageConfig
// Access to the other clients
ConfigProvider RestConfigProvider
}
func NewStorageOptions() *StorageOptions {
@@ -140,7 +149,7 @@ func (o *StorageOptions) ApplyTo(serverConfig *genericapiserver.RecommendedConfi
if err != nil {
return err
}
getter := apistore.NewRESTOptionsGetterForClient(unified, etcdOptions.StorageConfig)
getter := apistore.NewRESTOptionsGetterForClient(unified, etcdOptions.StorageConfig, o.ConfigProvider)
serverConfig.RESTOptionsGetter = getter
return nil
}
+9 -6
View File
@@ -95,11 +95,12 @@ type service struct {
storageStatus dualwrite.Service
kvStore kvstore.KVStore
pluginClient plugins.Client
datasources datasource.ScopedPluginDatasourceProvider
contextProvider datasource.PluginContextWrapper
pluginStore pluginstore.Store
unified resource.ResourceClient
pluginClient plugins.Client
datasources datasource.ScopedPluginDatasourceProvider
contextProvider datasource.PluginContextWrapper
pluginStore pluginstore.Store
unified resource.ResourceClient
restConfigProvider RestConfigProvider
buildHandlerChainFuncFromBuilders builder.BuildHandlerChainFuncFromBuilders
}
@@ -118,6 +119,7 @@ func ProvideService(
pluginStore pluginstore.Store,
storageStatus dualwrite.Service,
unified resource.ResourceClient,
restConfigProvider RestConfigProvider,
buildHandlerChainFuncFromBuilders builder.BuildHandlerChainFuncFromBuilders,
eventualRestConfigProvider *eventualRestConfigProvider,
) (*service, error) {
@@ -144,6 +146,7 @@ func ProvideService(
serverLockService: serverLockService,
storageStatus: storageStatus,
unified: unified,
restConfigProvider: restConfigProvider,
buildHandlerChainFuncFromBuilders: buildHandlerChainFuncFromBuilders,
}
// This will be used when running as a dskit service
@@ -300,7 +303,7 @@ func (s *service) start(ctx context.Context) error {
return err
}
} else {
getter := apistore.NewRESTOptionsGetterForClient(s.unified, o.RecommendedOptions.Etcd.StorageConfig)
getter := apistore.NewRESTOptionsGetterForClient(s.unified, o.RecommendedOptions.Etcd.StorageConfig, s.restConfigProvider)
optsregister = getter.RegisterOptions
// Use unified storage client
+88 -6
View File
@@ -1,15 +1,25 @@
package apistore
import (
"context"
"errors"
"fmt"
"net/http"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apiserver/pkg/storage"
"k8s.io/client-go/rest"
authtypes "github.com/grafana/authlib/types"
"github.com/grafana/grafana/pkg/apimachinery/utils"
provisioningV0 "github.com/grafana/grafana/pkg/generated/clientset/versioned/typed/provisioning/v0alpha1"
"github.com/grafana/grafana/pkg/storage/unified/resource"
)
var errResourceIsManagedInRepository = fmt.Errorf("this resource is managed by a repository")
func checkManagerPropertiesOnDelete(auth authtypes.AuthInfo, obj utils.GrafanaMetaAccessor) error {
return enforceManagerProperties(auth, obj)
}
@@ -65,12 +75,8 @@ func enforceManagerProperties(auth authtypes.AuthInfo, obj utils.GrafanaMetaAcce
if auth.GetUID() == "access-policy:provisioning" {
return nil // OK!
}
return &apierrors.StatusError{ErrStatus: metav1.Status{
Status: metav1.StatusFailure,
Code: http.StatusForbidden,
Reason: metav1.StatusReasonForbidden,
Message: "Provisioned resources must be manaaged by the provisioning service account",
}}
// This can fallback to writing the value with a provisioning client
return errResourceIsManagedInRepository
case utils.ManagerKindPlugin, utils.ManagerKindClassicFP: // nolint:staticcheck
// ?? what identity do we use for legacy internal requests?
@@ -86,3 +92,79 @@ func enforceManagerProperties(auth authtypes.AuthInfo, obj utils.GrafanaMetaAcce
}
return nil
}
func (s *Storage) handleManagedResourceRouting(ctx context.Context,
err error,
action resource.WatchEvent_Type,
key string,
orig runtime.Object,
rsp runtime.Object,
) error {
if !errors.Is(err, errResourceIsManagedInRepository) || s.configProvider == nil {
return err
}
obj, err := utils.MetaAccessor(orig)
if err != nil {
return err
}
repo, ok := obj.GetManagerProperties()
if !ok {
return fmt.Errorf("expected managed resource")
}
if repo.Kind != utils.ManagerKindRepo {
return fmt.Errorf("expected managed repository")
}
src, ok := obj.GetSourceProperties()
if !ok || src.Path == "" {
return fmt.Errorf("missing source properties")
}
cfg, err := s.configProvider.GetRestConfig(ctx)
if err != nil {
return err
}
client, err := provisioningV0.NewForConfig(cfg)
if err != nil {
return err
}
if action == resource.WatchEvent_DELETED {
// TODO? can we copy orig into rsp without a full get?
if err = s.Get(ctx, key, storage.GetOptions{}, rsp); err != nil { // COPY?
return err
}
result := client.RESTClient().Delete().
Namespace(obj.GetNamespace()).
Resource("repositories").
Name(repo.Identity).
Suffix("files", src.Path).
Do(ctx)
return result.Error()
}
var req *rest.Request
switch action {
case resource.WatchEvent_ADDED:
req = client.RESTClient().Post()
case resource.WatchEvent_MODIFIED:
req = client.RESTClient().Put()
default:
return fmt.Errorf("unsupported provisioning action: %v, %w", action, err)
}
// Execute the change
result := req.Namespace(obj.GetNamespace()).
Resource("repositories").
Name(repo.Identity).
Suffix("files", src.Path).
Body(orig).
Param("skipDryRun", "true").
Do(ctx)
err = result.Error()
if err != nil {
return err
}
// return the updated value
return s.Get(ctx, key, storage.GetOptions{}, rsp)
}
+11 -7
View File
@@ -27,18 +27,20 @@ var _ generic.RESTOptionsGetter = (*RESTOptionsGetter)(nil)
type StorageOptionsRegister func(gr schema.GroupResource, opts StorageOptions)
type RESTOptionsGetter struct {
client resource.ResourceClient
original storagebackend.Config
client resource.ResourceClient
original storagebackend.Config
configProvider RestConfigProvider
// Each group+resource may need custom options
options map[string]StorageOptions
}
func NewRESTOptionsGetterForClient(client resource.ResourceClient, original storagebackend.Config) *RESTOptionsGetter {
func NewRESTOptionsGetterForClient(client resource.ResourceClient, original storagebackend.Config, configProvider RestConfigProvider) *RESTOptionsGetter {
return &RESTOptionsGetter{
client: client,
original: original,
options: make(map[string]StorageOptions),
client: client,
original: original,
options: make(map[string]StorageOptions),
configProvider: configProvider,
}
}
@@ -58,6 +60,7 @@ func NewRESTOptionsGetterMemory(originalStorageConfig storagebackend.Config) (*R
return NewRESTOptionsGetterForClient(
resource.NewLocalResourceClient(server),
originalStorageConfig,
nil,
), nil
}
@@ -93,6 +96,7 @@ func NewRESTOptionsGetterForFile(path string,
return NewRESTOptionsGetterForClient(
resource.NewLocalResourceClient(server),
originalStorageConfig,
nil,
), nil
}
@@ -134,7 +138,7 @@ func (r *RESTOptionsGetter) GetRESTOptions(resource schema.GroupResource, _ runt
indexers *cache.Indexers,
) (storage.Interface, factory.DestroyFunc, error) {
return NewStorage(config, r.client, keyFunc, nil, newFunc, newListFunc, getAttrsFunc,
trigger, indexers, r.options[resource.String()])
trigger, indexers, r.configProvider, r.options[resource.String()])
},
DeleteCollectionWorkers: 0,
EnableGarbageCollection: false,
+24 -17
View File
@@ -16,6 +16,7 @@ import (
"strconv"
"time"
"github.com/bwmarrin/snowflake"
"golang.org/x/exp/rand"
apierrors "k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/api/meta"
@@ -27,10 +28,9 @@ import (
"k8s.io/apiserver/pkg/storage"
"k8s.io/apiserver/pkg/storage/storagebackend"
"k8s.io/apiserver/pkg/storage/storagebackend/factory"
clientrest "k8s.io/client-go/rest"
"k8s.io/client-go/tools/cache"
"github.com/bwmarrin/snowflake"
authtypes "github.com/grafana/authlib/types"
"github.com/grafana/grafana/pkg/apimachinery/utils"
grafanaregistry "github.com/grafana/grafana/pkg/apiserver/registry/generic"
@@ -75,9 +75,10 @@ type Storage struct {
trigger storage.IndexerFuncs
indexers *cache.Indexers
store resource.ResourceClient
getKey func(string) (*resource.ResourceKey, error)
snowflake *snowflake.Node // used to enforce internal ids
store resource.ResourceClient
getKey func(string) (*resource.ResourceKey, error)
snowflake *snowflake.Node // used to enforce internal ids
configProvider RestConfigProvider // used for provisioning
versioner storage.Versioner
@@ -91,6 +92,10 @@ var ErrFileNotExists = fmt.Errorf("file doesn't exist")
// ErrNamespaceNotExists means the directory for the namespace doesn't actually exist.
var ErrNamespaceNotExists = errors.New("namespace does not exist")
type RestConfigProvider interface {
GetRestConfig(context.Context) (*clientrest.Config, error)
}
// NewStorage instantiates a new Storage.
func NewStorage(
config *storagebackend.ConfigForResource,
@@ -102,18 +107,20 @@ func NewStorage(
getAttrsFunc storage.AttrFunc,
trigger storage.IndexerFuncs,
indexers *cache.Indexers,
configProvider RestConfigProvider,
opts StorageOptions,
) (storage.Interface, factory.DestroyFunc, error) {
s := &Storage{
store: store,
gr: config.GroupResource,
codec: config.Codec,
keyFunc: keyFunc,
newFunc: newFunc,
newListFunc: newListFunc,
getAttrsFunc: getAttrsFunc,
trigger: trigger,
indexers: indexers,
store: store,
gr: config.GroupResource,
codec: config.Codec,
keyFunc: keyFunc,
newFunc: newFunc,
newListFunc: newListFunc,
getAttrsFunc: getAttrsFunc,
trigger: trigger,
indexers: indexers,
configProvider: configProvider,
getKey: keyParser,
@@ -173,7 +180,7 @@ func (s *Storage) Create(ctx context.Context, key string, obj runtime.Object, ou
req := &resource.CreateRequest{}
req.Value, permissions, err = s.prepareObjectForStorage(ctx, obj)
if err != nil {
return err
return s.handleManagedResourceRouting(ctx, err, resource.WatchEvent_ADDED, key, obj, out)
}
req.Key, err = s.getKey(key)
@@ -280,7 +287,7 @@ func (s *Storage) Delete(
return fmt.Errorf("unable to read object %w", err)
}
if err = checkManagerPropertiesOnDelete(info, meta); err != nil {
return err
return s.handleManagedResourceRouting(ctx, err, resource.WatchEvent_DELETED, key, out, out)
}
rsp, err := s.store.Delete(ctx, cmd)
@@ -580,7 +587,7 @@ func (s *Storage) GuaranteedUpdate(
req.Value, err = s.prepareObjectForUpdate(ctx, updatedObj, existingObj)
if err != nil {
return err
return s.handleManagedResourceRouting(ctx, err, resource.WatchEvent_MODIFIED, key, updatedObj, destination)
}
var rv uint64
@@ -185,6 +185,7 @@ func TestGRPCtoHTTPStatusMapping(t *testing.T) {
nil,
nil,
nil,
nil,
apistore.StorageOptions{})
require.NoError(t, err)
+1 -1
View File
@@ -180,7 +180,7 @@ func testSetup(t testing.TB, opts ...setupOption) (context.Context, storage.Inte
setupOpts.newListFunc,
storage.DefaultNamespaceScopedAttr,
make(map[string]storage.IndexerFunc, 0),
nil,
nil, nil,
apistore.StorageOptions{},
)
if err != nil {
@@ -18,13 +18,13 @@ import (
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/runtime"
"github.com/grafana/grafana/pkg/apimachinery/utils"
provisioning "github.com/grafana/grafana/pkg/apis/provisioning/v0alpha1"
"github.com/grafana/grafana/pkg/infra/slugify"
"github.com/grafana/grafana/pkg/infra/usagestats"
"github.com/grafana/grafana/pkg/tests/apis"
"k8s.io/apimachinery/pkg/runtime"
)
func TestIntegrationProvisioning_CreatingAndGetting(t *testing.T) {
@@ -556,18 +556,14 @@ func TestIntegrationProvisioning_ImportAllPanelsFromLocalRepository(t *testing.T
// Try writing the value directly
err = unstructured.SetNestedField(obj.Object, []any{"aaa", "bbb"}, "spec", "tags")
require.NoError(t, err, "set tags")
_, err = helper.Dashboards.Resource.Update(ctx, obj, metav1.UpdateOptions{})
require.Error(t, err, "only the provisionding service should be able to update")
require.True(t, apierrors.IsForbidden(err))
obj, err = helper.Dashboards.Resource.Update(ctx, obj, metav1.UpdateOptions{})
require.NoError(t, err)
v, _, _ := unstructured.NestedString(obj.Object, "metadata", "annotations", utils.AnnoKeyUpdatedBy)
require.Equal(t, "access-policy:provisioning", v)
// Should not be able to directly delete the managed resource
err = helper.Dashboards.Resource.Delete(ctx, allPanels, metav1.DeleteOptions{})
require.Error(t, err, "only the provisioning service should be able to delete")
require.True(t, apierrors.IsForbidden(err))
// But we can delete the repository file, and this should also remove the resource
err = helper.Repositories.Resource.Delete(ctx, repo, metav1.DeleteOptions{}, "files", "all-panels.json")
require.NoError(t, err, "should delete the resource file")
require.NoError(t, err, "user can delete")
_, err = helper.Dashboards.Resource.Get(ctx, allPanels, metav1.GetOptions{})
require.Error(t, err, "should delete the internal resource")