K8s/Dashboards: Delegate large objects to blob store (#94943)
This commit is contained in:
@@ -0,0 +1,166 @@
|
||||
package apistore
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/apimachinery/pkg/runtime/schema"
|
||||
|
||||
common "github.com/grafana/grafana/pkg/apimachinery/apis/common/v0alpha1"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/utils"
|
||||
"github.com/grafana/grafana/pkg/storage/unified/resource"
|
||||
)
|
||||
|
||||
type LargeObjectSupport interface {
|
||||
// The resource this can process
|
||||
GroupResource() schema.GroupResource
|
||||
|
||||
// The size that triggers delegating part of the object to blob storage
|
||||
Threshold() int
|
||||
|
||||
// Each resource may have a maximum size that is different than the global maximum
|
||||
// for example, we know we will allow dashboards up to 10mb, however most
|
||||
// resources should have a smaller limit (1mb?)
|
||||
MaxSize() int
|
||||
|
||||
// Deconstruct takes a large object, write most of it to blob storage and leave a few metadata bits around to help with list
|
||||
// NOTE: changes to the object must be handled by mutating the input obj
|
||||
Deconstruct(ctx context.Context, key *resource.ResourceKey, client resource.BlobStoreClient, obj utils.GrafanaMetaAccessor, raw []byte) error
|
||||
|
||||
// Reconstruct will join the resource+blob back into a complete resource
|
||||
// NOTE: changes to the object must be handled by mutating the input obj
|
||||
Reconstruct(ctx context.Context, key *resource.ResourceKey, client resource.BlobStoreClient, obj utils.GrafanaMetaAccessor) error
|
||||
}
|
||||
|
||||
var _ LargeObjectSupport = (*BasicLargeObjectSupport)(nil)
|
||||
|
||||
type BasicLargeObjectSupport struct {
|
||||
TheGroupResource schema.GroupResource
|
||||
ThresholdSize int
|
||||
MaxByteSize int
|
||||
|
||||
// Mutate the spec so it only has the small properties
|
||||
ReduceSpec func(obj runtime.Object) error
|
||||
|
||||
// Update the spec so it has the full object
|
||||
// This is used to support server-side apply
|
||||
RebuildSpec func(obj runtime.Object, blob []byte) error
|
||||
}
|
||||
|
||||
func (s *BasicLargeObjectSupport) GroupResource() schema.GroupResource {
|
||||
return s.TheGroupResource
|
||||
}
|
||||
|
||||
// Threshold implements LargeObjectSupport.
|
||||
func (s *BasicLargeObjectSupport) Threshold() int {
|
||||
return s.ThresholdSize
|
||||
}
|
||||
|
||||
// MaxSize implements LargeObjectSupport.
|
||||
func (s *BasicLargeObjectSupport) MaxSize() int {
|
||||
return s.MaxByteSize
|
||||
}
|
||||
|
||||
// Deconstruct implements LargeObjectSupport.
|
||||
func (s *BasicLargeObjectSupport) Deconstruct(ctx context.Context, key *resource.ResourceKey, client resource.BlobStoreClient, obj utils.GrafanaMetaAccessor, raw []byte) error {
|
||||
if key.Group != s.TheGroupResource.Group {
|
||||
return fmt.Errorf("requested group mismatch")
|
||||
}
|
||||
if key.Resource != s.TheGroupResource.Resource {
|
||||
return fmt.Errorf("requested resource mismatch")
|
||||
}
|
||||
|
||||
spec, err := obj.GetSpec()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
var val []byte
|
||||
|
||||
// :( could not figure out custom JSON marshaling
|
||||
// with pointer receiver... this is a quick fix to support dashboards
|
||||
u, ok := spec.(common.Unstructured)
|
||||
if ok {
|
||||
val, err = json.Marshal(u.Object)
|
||||
} else {
|
||||
val, err = json.Marshal(spec)
|
||||
}
|
||||
|
||||
// Write only the spec
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
rt, ok := obj.GetRuntimeObject()
|
||||
if !ok {
|
||||
return fmt.Errorf("expected runtime object")
|
||||
}
|
||||
|
||||
err = s.ReduceSpec(rt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Save the blob
|
||||
info, err := client.PutBlob(ctx, &resource.PutBlobRequest{
|
||||
ContentType: "application/json",
|
||||
Value: val,
|
||||
Resource: key,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Update the resource metadata with the blob info
|
||||
obj.SetBlob(&utils.BlobInfo{
|
||||
UID: info.Uid,
|
||||
Size: info.Size,
|
||||
Hash: info.Hash,
|
||||
MimeType: info.MimeType,
|
||||
Charset: info.Charset,
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// Reconstruct implements LargeObjectSupport.
|
||||
func (s *BasicLargeObjectSupport) Reconstruct(ctx context.Context, key *resource.ResourceKey, client resource.BlobStoreClient, obj utils.GrafanaMetaAccessor) error {
|
||||
blobInfo := obj.GetBlob()
|
||||
if blobInfo == nil {
|
||||
return fmt.Errorf("the object does not have a blob")
|
||||
}
|
||||
|
||||
rv, err := obj.GetResourceVersionInt64()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
rsp, err := client.GetBlob(ctx, &resource.GetBlobRequest{
|
||||
Resource: &resource.ResourceKey{
|
||||
Group: s.TheGroupResource.Group,
|
||||
Resource: s.TheGroupResource.Resource,
|
||||
Namespace: obj.GetNamespace(),
|
||||
Name: obj.GetName(),
|
||||
},
|
||||
MustProxyBytes: true,
|
||||
ResourceVersion: rv,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if rsp.Error != nil {
|
||||
return fmt.Errorf("error loading value from object store %+v", rsp.Error)
|
||||
}
|
||||
|
||||
// Replace the spec with the value saved in the blob store
|
||||
if len(rsp.Value) == 0 {
|
||||
return fmt.Errorf("empty blob value")
|
||||
}
|
||||
|
||||
rt, ok := obj.GetRuntimeObject()
|
||||
if !ok {
|
||||
return fmt.Errorf("unable to get raw object")
|
||||
}
|
||||
obj.SetBlob(nil) // remove the blob info
|
||||
return s.RebuildSpec(rt, rsp.Value)
|
||||
}
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"math"
|
||||
"time"
|
||||
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
@@ -15,6 +16,23 @@ import (
|
||||
"github.com/grafana/grafana/pkg/storage/unified/resource"
|
||||
)
|
||||
|
||||
func logN(n, b float64) float64 {
|
||||
return math.Log(n) / math.Log(b)
|
||||
}
|
||||
|
||||
// Slightly modified function from https://github.com/dustin/go-humanize (MIT).
|
||||
func formatBytes(numBytes int) string {
|
||||
base := 1024.0
|
||||
sizes := []string{"B", "KiB", "MiB", "GiB", "TiB", "PiB", "EiB"}
|
||||
if numBytes < 10 {
|
||||
return fmt.Sprintf("%d B", numBytes)
|
||||
}
|
||||
e := math.Floor(logN(float64(numBytes), base))
|
||||
suffix := sizes[int(e)]
|
||||
val := math.Floor(float64(numBytes)/math.Pow(base, e)*10+0.5) / 10
|
||||
return fmt.Sprintf("%.1f %s", val, suffix)
|
||||
}
|
||||
|
||||
// Called on create
|
||||
func (s *Storage) prepareObjectForStorage(ctx context.Context, newObject runtime.Object) ([]byte, error) {
|
||||
user, err := identity.GetRequester(ctx)
|
||||
@@ -51,11 +69,7 @@ func (s *Storage) prepareObjectForStorage(ctx context.Context, newObject runtime
|
||||
if err = s.codec.Encode(newObject, &buf); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if s.largeObjectSupport {
|
||||
return s.handleLargeResources(ctx, obj, buf)
|
||||
}
|
||||
return buf.Bytes(), nil
|
||||
return s.handleLargeResources(ctx, obj, buf)
|
||||
}
|
||||
|
||||
// Called on update
|
||||
@@ -106,29 +120,41 @@ func (s *Storage) prepareObjectForUpdate(ctx context.Context, updateObject runti
|
||||
if err = s.codec.Encode(updateObject, &buf); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if s.largeObjectSupport {
|
||||
return s.handleLargeResources(ctx, obj, buf)
|
||||
}
|
||||
return buf.Bytes(), nil
|
||||
return s.handleLargeResources(ctx, obj, buf)
|
||||
}
|
||||
|
||||
func (s *Storage) handleLargeResources(ctx context.Context, obj utils.GrafanaMetaAccessor, buf bytes.Buffer) ([]byte, error) {
|
||||
if buf.Len() > 1000 {
|
||||
// !!! Currently just write the whole thing
|
||||
// in reality we may only want to write the spec....
|
||||
_, err := s.store.PutBlob(ctx, &resource.PutBlobRequest{
|
||||
ContentType: "application/json",
|
||||
Value: buf.Bytes(),
|
||||
Resource: &resource.ResourceKey{
|
||||
Group: s.gr.Group,
|
||||
Resource: s.gr.Resource,
|
||||
Namespace: obj.GetNamespace(),
|
||||
Name: obj.GetName(),
|
||||
},
|
||||
})
|
||||
support := s.opts.LargeObjectSupport
|
||||
if support != nil {
|
||||
size := buf.Len()
|
||||
if size > support.Threshold() {
|
||||
if support.MaxSize() > 0 && size > support.MaxSize() {
|
||||
return nil, fmt.Errorf("request object is too big (%s > %s)", formatBytes(size), formatBytes(support.MaxSize()))
|
||||
}
|
||||
}
|
||||
|
||||
key := &resource.ResourceKey{
|
||||
Group: s.gr.Group,
|
||||
Resource: s.gr.Resource,
|
||||
Namespace: obj.GetNamespace(),
|
||||
Name: obj.GetName(),
|
||||
}
|
||||
|
||||
err := support.Deconstruct(ctx, key, s.store, obj, buf.Bytes())
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
buf.Reset()
|
||||
orig, ok := obj.GetRuntimeObject()
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("error using object as runtime object")
|
||||
}
|
||||
|
||||
// Now encode the smaller version
|
||||
if err = s.codec.Encode(orig, &buf); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
return buf.Bytes(), nil
|
||||
}
|
||||
|
||||
@@ -24,25 +24,25 @@ import (
|
||||
|
||||
var _ generic.RESTOptionsGetter = (*RESTOptionsGetter)(nil)
|
||||
|
||||
// This is a copy of the original flag, as we are not allowed to import grafana core.
|
||||
const bigObjectSupportFlag = "unifiedStorageBigObjectsSupport"
|
||||
type StorageOptionsRegister func(gr schema.GroupResource, opts StorageOptions)
|
||||
|
||||
type RESTOptionsGetter struct {
|
||||
client resource.ResourceClient
|
||||
original storagebackend.Config
|
||||
// As we are not allowed to import the feature management directly, we pass a map of enabled features.
|
||||
features map[string]any
|
||||
|
||||
// Each group+resource may need custom options
|
||||
options map[string]StorageOptions
|
||||
}
|
||||
|
||||
func NewRESTOptionsGetterForClient(client resource.ResourceClient, original storagebackend.Config, features map[string]any) *RESTOptionsGetter {
|
||||
func NewRESTOptionsGetterForClient(client resource.ResourceClient, original storagebackend.Config) *RESTOptionsGetter {
|
||||
return &RESTOptionsGetter{
|
||||
client: client,
|
||||
original: original,
|
||||
features: features,
|
||||
options: make(map[string]StorageOptions),
|
||||
}
|
||||
}
|
||||
|
||||
func NewRESTOptionsGetterMemory(originalStorageConfig storagebackend.Config, features map[string]any) (*RESTOptionsGetter, error) {
|
||||
func NewRESTOptionsGetterMemory(originalStorageConfig storagebackend.Config) (*RESTOptionsGetter, error) {
|
||||
backend, err := resource.NewCDKBackend(context.Background(), resource.CDKBackendOptions{
|
||||
Bucket: memblob.OpenBucket(&memblob.Options{}),
|
||||
})
|
||||
@@ -58,7 +58,6 @@ func NewRESTOptionsGetterMemory(originalStorageConfig storagebackend.Config, fea
|
||||
return NewRESTOptionsGetterForClient(
|
||||
resource.NewLocalResourceClient(server),
|
||||
originalStorageConfig,
|
||||
features,
|
||||
), nil
|
||||
}
|
||||
|
||||
@@ -94,10 +93,13 @@ func NewRESTOptionsGetterForFile(path string,
|
||||
return NewRESTOptionsGetterForClient(
|
||||
resource.NewLocalResourceClient(server),
|
||||
originalStorageConfig,
|
||||
features,
|
||||
), nil
|
||||
}
|
||||
|
||||
func (r *RESTOptionsGetter) RegisterOptions(gr schema.GroupResource, opts StorageOptions) {
|
||||
r.options[gr.String()] = opts
|
||||
}
|
||||
|
||||
// TODO: The RESTOptionsGetter interface added a new example object parameter to help determine the default
|
||||
// storage version for a resource. This is not currently used in this implementation.
|
||||
func (r *RESTOptionsGetter) GetRESTOptions(resource schema.GroupResource, _ runtime.Object) (generic.RESTOptions, error) {
|
||||
@@ -131,12 +133,8 @@ func (r *RESTOptionsGetter) GetRESTOptions(resource schema.GroupResource, _ runt
|
||||
trigger storage.IndexerFuncs,
|
||||
indexers *cache.Indexers,
|
||||
) (storage.Interface, factory.DestroyFunc, error) {
|
||||
if _, enabled := r.features[bigObjectSupportFlag]; enabled {
|
||||
return NewStorage(config, r.client, keyFunc, nil, newFunc, newListFunc, getAttrsFunc,
|
||||
trigger, indexers, LargeObjectSupportEnabled)
|
||||
}
|
||||
return NewStorage(config, r.client, keyFunc, nil, newFunc, newListFunc, getAttrsFunc,
|
||||
trigger, indexers, LargeObjectSupportDisabled)
|
||||
trigger, indexers, r.options[resource.String()])
|
||||
},
|
||||
DeleteCollectionWorkers: 0,
|
||||
EnableGarbageCollection: false,
|
||||
|
||||
@@ -41,6 +41,11 @@ const (
|
||||
|
||||
var _ storage.Interface = (*Storage)(nil)
|
||||
|
||||
// Optional settings that apply to a single resource
|
||||
type StorageOptions struct {
|
||||
LargeObjectSupport LargeObjectSupport
|
||||
}
|
||||
|
||||
// Storage implements storage.Interface and storage resources as JSON files on disk.
|
||||
type Storage struct {
|
||||
gr schema.GroupResource
|
||||
@@ -57,9 +62,8 @@ type Storage struct {
|
||||
|
||||
versioner storage.Versioner
|
||||
|
||||
// Defines if we want to outsource large objects to another storage type.
|
||||
// By default, this feature is disabled.
|
||||
largeObjectSupport bool
|
||||
// Resource options like large object support
|
||||
opts StorageOptions
|
||||
}
|
||||
|
||||
// ErrFileNotExists means the file doesn't actually exist.
|
||||
@@ -79,7 +83,7 @@ func NewStorage(
|
||||
getAttrsFunc storage.AttrFunc,
|
||||
trigger storage.IndexerFuncs,
|
||||
indexers *cache.Indexers,
|
||||
largeObjectSupport bool,
|
||||
opts StorageOptions,
|
||||
) (storage.Interface, factory.DestroyFunc, error) {
|
||||
s := &Storage{
|
||||
store: store,
|
||||
@@ -96,7 +100,7 @@ func NewStorage(
|
||||
|
||||
versioner: &storage.APIObjectVersioner{},
|
||||
|
||||
largeObjectSupport: largeObjectSupport,
|
||||
opts: opts,
|
||||
}
|
||||
|
||||
// The key parsing callback allows us to support the hardcoded paths from upstream tests
|
||||
@@ -480,6 +484,14 @@ func (s *Storage) GuaranteedUpdate(
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
// restore the full original object before tryUpdate
|
||||
if s.opts.LargeObjectSupport != nil && mmm.GetBlob() != nil {
|
||||
err = s.opts.LargeObjectSupport.Reconstruct(ctx, req.Key, s.store, mmm)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
} else if !ignoreNotFound {
|
||||
return apierrors.NewNotFound(s.gr, req.Key.Name)
|
||||
}
|
||||
|
||||
@@ -176,7 +176,7 @@ func testSetup(t testing.TB, opts ...setupOption) (context.Context, storage.Inte
|
||||
storage.DefaultNamespaceScopedAttr,
|
||||
make(map[string]storage.IndexerFunc, 0),
|
||||
nil,
|
||||
LargeObjectSupportDisabled,
|
||||
StorageOptions{},
|
||||
)
|
||||
if err != nil {
|
||||
return nil, nil, nil, err
|
||||
|
||||
@@ -7,14 +7,16 @@ import (
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"mime"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/utils"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
"go.opentelemetry.io/otel/trace/noop"
|
||||
"gocloud.dev/blob"
|
||||
|
||||
"github.com/grafana/grafana/pkg/apimachinery/utils"
|
||||
|
||||
// Supported drivers
|
||||
_ "gocloud.dev/blob/azureblob"
|
||||
_ "gocloud.dev/blob/fileblob"
|
||||
@@ -32,6 +34,14 @@ type CDKBlobSupportOptions struct {
|
||||
|
||||
// Called in a context that loaded the possible drivers
|
||||
func OpenBlobBucket(ctx context.Context, url string) (*blob.Bucket, error) {
|
||||
if strings.HasPrefix(url, "file:") {
|
||||
// Don't write metadata attributes
|
||||
if strings.Contains(url, "?") {
|
||||
url += "&metadata=skip"
|
||||
} else {
|
||||
url += "?metadata=skip"
|
||||
}
|
||||
}
|
||||
return blob.OpenBucket(ctx, url)
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,18 @@
|
||||
package resource
|
||||
|
||||
func verifyRequestKey(key *ResourceKey) *ErrorResult {
|
||||
if key == nil {
|
||||
return NewBadRequestError("missing resource key")
|
||||
}
|
||||
if key.Group == "" {
|
||||
return NewBadRequestError("request key is missing group")
|
||||
}
|
||||
if key.Resource == "" {
|
||||
return NewBadRequestError("request key is missing resource")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func matchesQueryKey(query *ResourceKey, key *ResourceKey) bool {
|
||||
if query.Group != key.Group {
|
||||
return false
|
||||
|
||||
@@ -10,15 +10,16 @@ import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/authlib/claims"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/identity"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/utils"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
"go.opentelemetry.io/otel/trace/noop"
|
||||
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
||||
|
||||
"github.com/grafana/authlib/claims"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/identity"
|
||||
"github.com/grafana/grafana/pkg/apimachinery/utils"
|
||||
)
|
||||
|
||||
// ResourceServer implements all gRPC services
|
||||
@@ -821,6 +822,10 @@ func (s *server) PutBlob(ctx context.Context, req *PutBlobRequest) (*PutBlobResp
|
||||
}
|
||||
|
||||
func (s *server) getPartialObject(ctx context.Context, key *ResourceKey, rv int64) (utils.GrafanaMetaAccessor, *ErrorResult) {
|
||||
if r := verifyRequestKey(key); r != nil {
|
||||
return nil, r
|
||||
}
|
||||
|
||||
rsp := s.backend.ReadResource(ctx, &ReadRequest{
|
||||
Key: key,
|
||||
ResourceVersion: rv,
|
||||
|
||||
Reference in New Issue
Block a user