From 3a9559c8d06109c20661157bb9175b0b80dd443c Mon Sep 17 00:00:00 2001 From: konsalex Date: Wed, 26 Nov 2025 11:41:50 +0100 Subject: [PATCH] rebasing --- conf/apiextensions.ini | 48 ++ go.work.sum | 9 + pkg/apiserver/registry/generic/key.go | 40 +- pkg/apiserver/registry/generic/storage.go | 19 +- pkg/extensions/enterprise_imports.go | 6 +- pkg/registry/apis/apiextensions/README.md | 113 +++ pkg/registry/apis/apiextensions/discovery.go | 167 +++++ .../apis/apiextensions/dynamic_handler.go | 451 ++++++++++++ .../apis/apiextensions/dynamic_registry.go | 452 ++++++++++++ pkg/registry/apis/apiextensions/register.go | 362 +++++++++ .../apiextensions/resources/example-crd.yaml | 69 ++ .../resources/example-widget.yaml | 9 + .../apiextensions/resources/outage-crd.yaml | 111 +++ .../resources/outage-instance.yaml | 22 + pkg/registry/apis/apiextensions/storage.go | 176 +++++ pkg/registry/apis/apiextensions/storage_cr.go | 687 ++++++++++++++++++ pkg/registry/apis/apis.go | 2 + pkg/registry/apis/wireset.go | 2 + pkg/server/wire_gen.go | 13 +- pkg/services/apiserver/service.go | 86 ++- pkg/services/featuremgmt/registry.go | 7 + pkg/services/featuremgmt/toggles_gen.go | 4 + pkg/storage/unified/resource/access.go | 6 +- 23 files changed, 2836 insertions(+), 25 deletions(-) create mode 100644 conf/apiextensions.ini create mode 100644 pkg/registry/apis/apiextensions/README.md create mode 100644 pkg/registry/apis/apiextensions/discovery.go create mode 100644 pkg/registry/apis/apiextensions/dynamic_handler.go create mode 100644 pkg/registry/apis/apiextensions/dynamic_registry.go create mode 100644 pkg/registry/apis/apiextensions/register.go create mode 100644 pkg/registry/apis/apiextensions/resources/example-crd.yaml create mode 100644 pkg/registry/apis/apiextensions/resources/example-widget.yaml create mode 100644 pkg/registry/apis/apiextensions/resources/outage-crd.yaml create mode 100644 pkg/registry/apis/apiextensions/resources/outage-instance.yaml create mode 100644 pkg/registry/apis/apiextensions/storage.go create mode 100644 pkg/registry/apis/apiextensions/storage_cr.go diff --git a/conf/apiextensions.ini b/conf/apiextensions.ini new file mode 100644 index 00000000000..39884b5e656 --- /dev/null +++ b/conf/apiextensions.ini @@ -0,0 +1,48 @@ +; Run locally unified storage with SQLite to test +; new API registration changes +app_mode = development +target = all + +[log] +level = debug + +[server] +; HTTPS is required for kubectl (but HTTP works for testing with curl) +protocol = https +http_port = 1111 + +[feature_toggles] +; Enable the apiextensions feature +apiExtensions = true +; Enable unified storage globally +unifiedStorage = true +; Enable search indexing for unified storage +unifiedStorageSearch = true +; Enable the grafana-apiserver explicitly +grafanaAPIServer = true + +[grafana-apiserver] +; Use unified storage backed by SQL (uses your Grafana database) +storage_type = unified + +; Configure dashboards to use unified storage +[unified_storage.dashboards.dashboard.grafana.app] +; Dualwriter modes: +; 0: disabled (default) - dashboards saved to SQL only +; 1: read from legacy, write to legacy, write to unified best-effort +; 2: read from legacy, write to both +; 3: read from unified, write to both +; 4: read from unified, write to unified (fully migrated) +; 5: read from unified, write to unified, ignore background sync state +dualWriterMode = 5 + +; Configure folders to use unified storage (required for dashboards) +[unified_storage.folders.folder.grafana.app] +dualWriterMode = 5 + +[database] +; SQLite database for testing +type = sqlite3 +path = grafana.db +; Enable high availability mode is false for single instance +high_availability = false \ No newline at end of file diff --git a/go.work.sum b/go.work.sum index c4b732eaadb..256eac9e3a5 100644 --- a/go.work.sum +++ b/go.work.sum @@ -407,6 +407,7 @@ github.com/apache/arrow/go/v15 v15.0.2/go.mod h1:DGXsR3ajT524njufqf95822i+KTh+ye github.com/apache/thrift v0.21.0/go.mod h1:W1H8aR/QRtYNvrPeFXBtobyRkd0/YVhTc6i07XIAgDw= github.com/armon/circbuf v0.0.0-20150827004946-bbbad097214e h1:QEF07wC0T1rKkctt1RINW/+RMTVmiwxETico2l3gxJA= github.com/armon/consul-api v0.0.0-20180202201655-eb2c6b5be1b6 h1:G1bPvciwNyF7IUmKXNt9Ak3m6u9DE1rF+RmtIkBpVdA= +github.com/at-wat/mqtt-go v0.19.4/go.mod h1:AsiWc9kqVOhqq7LzUeWT/AkKUBfx3Sw5cEe8lc06fqA= github.com/atc0005/go-teams-notify/v2 v2.13.0 h1:nbDeHy89NjYlF/PEfLVF6lsserY9O5SnN1iOIw3AxXw= github.com/atc0005/go-teams-notify/v2 v2.13.0/go.mod h1:WSv9moolRsBcpZbwEf6gZxj7h0uJlJskJq5zkEWKO8Y= github.com/atomicgo/cursor v0.0.1/go.mod h1:cBON2QmmrysudxNBFthvMtN32r3jxVRIvzkUiF/RuIk= @@ -846,8 +847,10 @@ github.com/gorilla/handlers v1.5.2/go.mod h1:dX+xVpaxdSw+q0Qek8SSsl3dfMk3jNddUkM github.com/gorilla/mux v1.8.0/go.mod h1:DVbg23sWSpFRCP0SfiEN6jmj59UnW/n46BH5rLB71So= github.com/gorilla/websocket v1.4.2/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= github.com/grafana/alerting v0.0.0-20250729175202-b4b881b7b263/go.mod h1:VKxaR93Gff0ZlO2sPcdPVob1a/UzArFEW5zx3Bpyhls= +github.com/grafana/alerting v0.0.0-20251009192429-9427c24835ae/go.mod h1:VGjS5gDwWEADPP6pF/drqLxEImgeuHlEW5u8E5EfIrM= github.com/grafana/authlib v0.0.0-20250710201142-9542f2f28d43/go.mod h1:1fWkOiL+m32NBgRHZtlZGz2ji868tPZACYbqP3nBRJI= github.com/grafana/authlib/types v0.0.0-20250710201142-9542f2f28d43/go.mod h1:qeWYbnWzaYGl88JlL9+DsP1GT2Cudm58rLtx13fKZdw= +github.com/grafana/authlib/types v0.0.0-20250926065801-df98203cff37/go.mod h1:qeWYbnWzaYGl88JlL9+DsP1GT2Cudm58rLtx13fKZdw= github.com/grafana/cloudflare-go v0.0.0-20230110200409-c627cf6792f2 h1:qhugDMdQ4Vp68H0tp/0iN17DM2ehRo1rLEdOFe/gB8I= github.com/grafana/cloudflare-go v0.0.0-20230110200409-c627cf6792f2/go.mod h1:w/aiO1POVIeXUQyl0VQSZjl5OAGDTL5aX+4v0RA1tcw= github.com/grafana/cog v0.0.43/go.mod h1:TDunc7TYF7EfzjwFOlC5AkMe3To/U2KqyyG3QVvrF38= @@ -896,6 +899,7 @@ github.com/grafana/grafana-plugin-sdk-go v0.277.0/go.mod h1:mAUWg68w5+1f5TLDqagI github.com/grafana/grafana-plugin-sdk-go v0.278.0/go.mod h1:+8NXT/XUJ/89GV6FxGQ366NZ3nU+cAXDMd0OUESF9H4= github.com/grafana/grafana-plugin-sdk-go v0.279.0/go.mod h1:/7oGN6Z7DGTGaLHhgIYrRr6Wvmdsb3BLw5hL4Kbjy88= github.com/grafana/grafana-plugin-sdk-go v0.280.0/go.mod h1:Z15Wiq3c4I0tzHYrLYpOqrO8u3+2RJ+HN2Q9uiZTILA= +github.com/grafana/grafana-plugin-sdk-go v0.281.0/go.mod h1:3I0g+v6jAwVmrt6BEjDUP4V6pkhGP5QKY5NkXY4Ayr4= github.com/grafana/grafana-plugin-sdk-go v0.283.0/go.mod h1:20qhoYxIgbZRmwCEO1KMP8q2yq/Kge5+xE/99/hLEk0= github.com/grafana/grafana/apps/advisor v0.0.0-20250123151950-b066a6313173/go.mod h1:goSDiy3jtC2cp8wjpPZdUHRENcoSUHae1/Px/MDfddA= github.com/grafana/grafana/apps/advisor v0.0.0-20250220154326-6e5de80ef295/go.mod h1:9I1dKV3Dqr0NPR9Af0WJGxOytp5/6W3JLiNChOz8r+c= @@ -923,6 +927,7 @@ github.com/grafana/nanogit v0.0.0-20250616082354-5e94194d02ed/go.mod h1:OIAAKNgG github.com/grafana/nanogit v0.0.0-20250619160700-ebf70d342aa5 h1:MAQ2B0cu0V1S91ZjVa7NomNZFjaR2SmdtvdwhqBtyhU= github.com/grafana/nanogit v0.0.0-20250619160700-ebf70d342aa5/go.mod h1:tN93IZUaAmnSWgL0IgnKdLv6DNeIhTJGvl1wvQMrWco= github.com/grafana/nanogit v0.0.0-20250723104447-68f58f5ecec0/go.mod h1:ToqLjIdvV3AZQa3K6e5m9hy/nsGaUByc2dWQlctB9iA= +github.com/grafana/nanogit v0.0.0-20251106115617-c622d3e0fc4b/go.mod h1:ToqLjIdvV3AZQa3K6e5m9hy/nsGaUByc2dWQlctB9iA= github.com/grafana/prometheus-alertmanager v0.25.1-0.20240930132144-b5e64e81e8d3 h1:6D2gGAwyQBElSrp3E+9lSr7k8gLuP3Aiy20rweLWeBw= github.com/grafana/prometheus-alertmanager v0.25.1-0.20240930132144-b5e64e81e8d3/go.mod h1:YeND+6FDA7OuFgDzYODN8kfPhXLCehcpxe4T9mdnpCY= github.com/grafana/prometheus-alertmanager v0.25.1-0.20250331083058-4563aec7a975 h1:4/BZkGObFWZf4cLbE2Vqg/1VTz67Q0AJ7LHspWLKJoQ= @@ -939,6 +944,7 @@ github.com/grpc-ecosystem/go-grpc-middleware v1.3.0/go.mod h1:z0ButlSOZa5vEBq9m2 github.com/grpc-ecosystem/go-grpc-middleware/providers/prometheus v1.0.1/go.mod h1:lXGCsh6c22WGtjr+qGHj1otzZpV/1kwTMAqkwZsnWRU= github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.1.0/go.mod h1:XKMd7iuf/RGPSMJ/U4HP0zS2Z9Fh8Ps9a+6X26m/tmI= github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.0/go.mod h1:qOchhhIlmRcqk/O9uCo/puJlyo07YINaIqdZfZG3Jkc= +github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.3.2/go.mod h1:wd1YpapPLivG6nQgbf7ZkG1hhSOXDhhn4MLTknx2aAc= github.com/grpc-ecosystem/grpc-gateway v1.16.0 h1:gmcG1KaJ57LophUzW0Hy8NmPhnMZb4M0+kPpLofRdBo= github.com/grpc-ecosystem/grpc-gateway/v2 v2.16.0/go.mod h1:YN5jB8ie0yfIUg6VvR9Kz84aCaG7AsGZnLjhHbUqwPg= github.com/grpc-ecosystem/grpc-gateway/v2 v2.26.3/go.mod h1:ndYquD05frm2vACXE1nsccT4oJzjhw2arTS2cpUD1PI= @@ -1324,6 +1330,7 @@ github.com/prometheus/common v0.62.0/go.mod h1:vyBcEuLSvWos9B1+CyL7JZ2up+uFzXhkq github.com/prometheus/common v0.64.0/go.mod h1:0gZns+BLRQ3V6NdaerOhMbwwRbNh9hkGINtQAsP5GS8= github.com/prometheus/common v0.65.0/go.mod h1:0gZns+BLRQ3V6NdaerOhMbwwRbNh9hkGINtQAsP5GS8= github.com/prometheus/common v0.66.1/go.mod h1:gcaUsgf3KfRSwHY4dIMXLPV0K/Wg1oZ8+SbZk/HH/dA= +github.com/prometheus/common v0.67.1/go.mod h1:RpmT9v35q2Y+lsieQsdOh5sXZ6ajUGC8NjZAmr8vb0Q= github.com/prometheus/common v0.67.2/go.mod h1:63W3KZb1JOKgcjlIr64WW/LvFGAqKPj0atm+knVGEko= github.com/prometheus/common/assets v0.2.0 h1:0P5OrzoHrYBOSM1OigWL3mY8ZvV2N4zIE/5AahrSrfM= github.com/prometheus/exporter-toolkit v0.10.1-0.20230714054209-2f4150c63f97/go.mod h1:LoBCZeRh+5hX+fSULNyFnagYlQG/gBsyA/deNzROkq8= @@ -1375,6 +1382,7 @@ github.com/schollz/closestmatch v2.1.0+incompatible h1:Uel2GXEpJqOWBrlyI+oY9LTiy github.com/schollz/closestmatch v2.1.0+incompatible/go.mod h1:RtP1ddjLong6gTkbtmuhtR2uUrrJOpYzYRvbcPAid+g= github.com/schollz/progressbar/v3 v3.14.6 h1:GyjwcWBAf+GFDMLziwerKvpuS7ZF+mNTAXIB2aspiZs= github.com/schollz/progressbar/v3 v3.14.6/go.mod h1:Nrzpuw3Nl0srLY0VlTvC4V6RL50pcEymjy6qyJAaLa0= +github.com/sclevine/spec v1.4.0 h1:z/Q9idDcay5m5irkZ28M7PtQM4aOISzOpj4bUPkDee8= github.com/sclevine/spec v1.4.0/go.mod h1:LvpgJaFyvQzRvc1kaDs0bulYwzC70PbiYjC4QnFHkOM= github.com/segmentio/fasthash v1.0.3 h1:EI9+KE1EwvMLBWwjpRDc+fEM+prwxDYbslddQGtrmhM= github.com/segmentio/fasthash v1.0.3/go.mod h1:waKX8l2N8yckOgmSsXJi7x1ZfdKZ4x7KRMzBtS3oedY= @@ -1932,6 +1940,7 @@ golang.org/x/sync v0.13.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= golang.org/x/sync v0.14.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= golang.org/x/sync v0.15.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= golang.org/x/sync v0.16.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= +golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= golang.org/x/sys v0.0.0-20210112080510-489259a85091/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210503080704-8803ae5d1324/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210616045830-e2b7044e8c71/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= diff --git a/pkg/apiserver/registry/generic/key.go b/pkg/apiserver/registry/generic/key.go index 01ce8eb89ea..d12790df9f9 100644 --- a/pkg/apiserver/registry/generic/key.go +++ b/pkg/apiserver/registry/generic/key.go @@ -131,19 +131,31 @@ func NamespaceKeyFunc(gr schema.GroupResource) func(ctx context.Context, name st } } -// NoNamespaceKeyFunc is the default function for constructing storage paths -// to a resource relative to the given prefix without a namespace. -func NoNamespaceKeyFunc(ctx context.Context, prefix string, gr schema.GroupResource, name string) (string, error) { - if len(name) == 0 { - return "", apierrors.NewBadRequest("Name parameter required.") +// ClusterScopedKeyFunc constructs storage paths for cluster-scoped resources (no namespace). +func ClusterScopedKeyFunc(gr schema.GroupResource) func(ctx context.Context, name string) (string, error) { + return func(ctx context.Context, name string) (string, error) { + if len(name) == 0 { + return "", apierrors.NewBadRequest("Name parameter required.") + } + if msgs := path.IsValidPathSegmentName(name); len(msgs) != 0 { + return "", apierrors.NewBadRequest(fmt.Sprintf("Name parameter invalid: %q: %s", name, strings.Join(msgs, ";"))) + } + key := &Key{ + Group: gr.Group, + Resource: gr.Resource, + Name: name, + } + return key.String(), nil + } +} + +// ClusterScopedKeyRootFunc is used by the generic registry store for cluster-scoped resources. +func ClusterScopedKeyRootFunc(gr schema.GroupResource) func(ctx context.Context) string { + return func(ctx context.Context) string { + key := &Key{ + Group: gr.Group, + Resource: gr.Resource, + } + return key.String() } - if msgs := path.IsValidPathSegmentName(name); len(msgs) != 0 { - return "", apierrors.NewBadRequest(fmt.Sprintf("Name parameter invalid: %q: %s", name, strings.Join(msgs, ";"))) - } - key := &Key{ - Group: gr.Group, - Resource: gr.Resource, - Name: name, - } - return prefix + key.String(), nil } diff --git a/pkg/apiserver/registry/generic/storage.go b/pkg/apiserver/registry/generic/storage.go index 98e2f1fe9df..953b5dfd56f 100644 --- a/pkg/apiserver/registry/generic/storage.go +++ b/pkg/apiserver/registry/generic/storage.go @@ -1,6 +1,8 @@ package generic import ( + "context" + "k8s.io/apimachinery/pkg/runtime" "k8s.io/apiserver/pkg/registry/generic" "k8s.io/apiserver/pkg/registry/generic/registry" @@ -12,16 +14,27 @@ func NewRegistryStore(scheme *runtime.Scheme, resourceInfo utils.ResourceInfo, o gv := resourceInfo.GroupVersion() gv.Version = runtime.APIVersionInternal strategy := NewStrategy(scheme, gv) + + gr := resourceInfo.GroupResource() + var keyRootFunc func(ctx context.Context) string + var keyFunc func(ctx context.Context, name string) (string, error) + if resourceInfo.IsClusterScoped() { strategy = strategy.WithClusterScope() + keyRootFunc = ClusterScopedKeyRootFunc(gr) + keyFunc = ClusterScopedKeyFunc(gr) + } else { + keyRootFunc = KeyRootFunc(gr) + keyFunc = NamespaceKeyFunc(gr) } + store := ®istry.Store{ NewFunc: resourceInfo.NewFunc, NewListFunc: resourceInfo.NewListFunc, - KeyRootFunc: KeyRootFunc(resourceInfo.GroupResource()), - KeyFunc: NamespaceKeyFunc(resourceInfo.GroupResource()), + KeyRootFunc: keyRootFunc, + KeyFunc: keyFunc, PredicateFunc: Matcher, - DefaultQualifiedResource: resourceInfo.GroupResource(), + DefaultQualifiedResource: gr, SingularQualifiedResource: resourceInfo.SingularGroupResource(), TableConvertor: resourceInfo.TableConverter(), CreateStrategy: strategy, diff --git a/pkg/extensions/enterprise_imports.go b/pkg/extensions/enterprise_imports.go index 472652cc103..113c2f8e4bb 100644 --- a/pkg/extensions/enterprise_imports.go +++ b/pkg/extensions/enterprise_imports.go @@ -15,7 +15,6 @@ import ( _ "github.com/blugelabs/bluge" _ "github.com/blugelabs/bluge_segment_api" _ "github.com/crewjam/saml" - _ "github.com/docker/go-connections/nat" _ "github.com/go-jose/go-jose/v4" _ "github.com/gobwas/glob" _ "github.com/googleapis/gax-go/v2" @@ -31,7 +30,6 @@ import ( _ "github.com/spf13/cobra" // used by the standalone apiserver cli _ "github.com/spyzhov/ajson" _ "github.com/stretchr/testify/require" - _ "github.com/testcontainers/testcontainers-go" _ "gocloud.dev/secrets/awskms" _ "gocloud.dev/secrets/azurekeyvault" _ "gocloud.dev/secrets/gcpkms" @@ -56,7 +54,9 @@ import ( _ "github.com/grafana/e2e" _ "github.com/grafana/gofpdf" _ "github.com/grafana/gomemcache/memcache" + _ "github.com/grafana/tempo/pkg/traceql" + _ "github.com/grafana/grafana/apps/alerting/alertenrichment/pkg/apis/alertenrichment/v1beta1" _ "github.com/grafana/grafana/apps/scope/pkg/apis/scope/v0alpha1" - _ "github.com/grafana/tempo/pkg/traceql" + _ "github.com/testcontainers/testcontainers-go" ) diff --git a/pkg/registry/apis/apiextensions/README.md b/pkg/registry/apis/apiextensions/README.md new file mode 100644 index 00000000000..6a2f3d5a9a1 --- /dev/null +++ b/pkg/registry/apis/apiextensions/README.md @@ -0,0 +1,113 @@ +# Grafana CRD Support - POC Testing Guide + +This directory contains example CRD definitions and custom resources for testing the Kubernetes CustomResourceDefinition (CRD) support in Grafana server. + + +## How to run + +To run just compile Grafana with `make build-go` and then run the service with the `apiextensions.ini` provided under the `conf` folder. + +```bash +./bin/darwin-arm64/grafana server --config conf/apiextensions.ini +``` + +To enable this feature we use the `apiExtensions = true` flag and also have unified storage as our storage backend. + +There is no need for US service to run in ST Grafana, and the database will be a SQLite one. + +## Testing Steps + +### Step 1: Create a CustomResourceDefinition + +Create the example CRD that defines a "Widget" resource: + +```bash +kubectl apply -f ./pkg/registry/apis/apiextensions/resources/example-crd.yaml +``` + +Or use curl after you set a Grafana Service account token: +```bash +export AUTH_SVC="Authorization: Bearer glsa_" +``` + +```bash +# From Grafana root +curl -k -X POST https://localhost:1111/apis/apiextensions.k8s.io/v1/customresourcedefinitions \ + -H "$AUTH_SVC" \ + -H "Content-Type: application/yaml" \ + --data-binary @$PWD/pkg/registry/apis/apiextensions/resources/example-crd.yaml +``` + +### Step 2: Verify the CRD was created + +List all CRDs: + +```bash +# Be sure to use the generated kube-config file +KUBECONFIG=$PWD/data/grafana-apiserver/apiserver.kubeconfig \ +kubectl get customresourcedefinitions.apiextensions.k8s.io + +# Or with curl: +curl -k -X GET https://localhost:1111/apis/apiextensions.k8s.io/v1/customresourcedefinitions \ + -H "$AUTH_SVC" +``` + +### Step 3: Create a Custom Resource Instance + +Now that the CRD is registered, create an instance of the Widget resource: + +```bash +kubectl apply -f ./pkg/registry/apis/apiextensions/resources/example-widget.yaml +``` + +Or use curl: + +```bash +curl -k -X POST https://localhost:1111/apis/apiextensions.k8s.io/v1/customresourcedefinitions \ + -H "$AUTH_SVC" \ + -H "Content-Type: application/yaml" \ + --data-binary @$PWD/pkg/registry/apis/apiextensions/resources/example-widget.yaml +``` + +### Step 4: Verify the Custom Resource + +> Note: You can see all the resources in SQLite database called `grafana.db` + +List all widgets: + +```bash +KUBECONFIG=$PWD/data/grafana-apiserver/apiserver.kubeconfig \ +kubectl get widgets -n default + +# Or with curl: +curl -k -X GET https://localhost:1111/apis/customcrdtest.grafana.app/v1/namespaces/default/widgets \ + -H "$AUTH_SVC" | jq . +``` + +### Step 5: Update the Custom Resource + +Update the widget's spec: + +```bash +KUBECONFIG=$PWD/data/grafana-apiserver/apiserver.kubeconfig \ +kubectl edit widget my-widget -n default + +# Or with curl (PATCH): +curl -X PATCH https://localhost:1111/apis/customcrdtest.grafana.app/v1/namespaces/default/widgets/my-widget \ + -H "Content-Type: application/merge-patch+json" \ + -H "$AUTH_SVC" \ + -d '{"spec":{"replicas":5}}' +``` + + + +## What is left + +- [ ] Support multiple versions of CRDs (example in `discovery.go`) +- [ ] Watch new CRDs, so we do not require server restart `dynamic_registry.go`. This is needed for horizontal deployments. +- [ ] Support `/status` subresource +- [ ] Add tracer and logger and remove `fmt.Print` +- [ ] Implement MT setup +- [ ] Figure out how to modify storage checks for Cluster scoped resources (when we create a new CRD) +- [ ] How to tackle Cluster scoped CRs +- [ ] Use the feature flag to start the `apiextensions` service on-demand diff --git a/pkg/registry/apis/apiextensions/discovery.go b/pkg/registry/apis/apiextensions/discovery.go new file mode 100644 index 00000000000..038b743322a --- /dev/null +++ b/pkg/registry/apis/apiextensions/discovery.go @@ -0,0 +1,167 @@ +package apiextensions + +import ( + "encoding/json" + "fmt" + "net/http" + "sync" + + apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime/schema" +) + +// DiscoveryManager manages discovery for dynamically registered custom resources +type DiscoveryManager struct { + mu sync.RWMutex + apiGroups map[string]*metav1.APIGroup // group name -> APIGroup + resources map[string]*metav1.APIResourceList // group/version -> APIResourceList +} + +// NewDiscoveryManager creates a new discovery manager +func NewDiscoveryManager() *DiscoveryManager { + return &DiscoveryManager{ + apiGroups: make(map[string]*metav1.APIGroup), + resources: make(map[string]*metav1.APIResourceList), + } +} + +// AddCustomResource adds a custom resource to the discovery documents +func (d *DiscoveryManager) AddCustomResource(crd *apiextensionsv1.CustomResourceDefinition) { + d.mu.Lock() + defer d.mu.Unlock() + + group := crd.Spec.Group + // TODO(@konsalex): is this needed to be iterated + // or multiple versions will be different crd entries? + // Need to test this out + version := crd.Spec.Versions[0].Name + gvKey := fmt.Sprintf("%s/%s", group, version) + + // Update API Group discovery + apiGroup, ok := d.apiGroups[group] + if !ok { + apiGroup = &metav1.APIGroup{ + TypeMeta: metav1.TypeMeta{ + Kind: "APIGroup", + APIVersion: "v1", + }, + Name: group, + Versions: []metav1.GroupVersionForDiscovery{ + { + GroupVersion: fmt.Sprintf("%s/%s", group, version), + Version: version, + }, + }, + PreferredVersion: metav1.GroupVersionForDiscovery{ + GroupVersion: fmt.Sprintf("%s/%s", group, version), + Version: version, + }, + } + d.apiGroups[group] = apiGroup + } + + // Update API Resource List + resourceList, ok := d.resources[gvKey] + if !ok { + resourceList = &metav1.APIResourceList{ + TypeMeta: metav1.TypeMeta{ + Kind: "APIResourceList", + APIVersion: "v1", + }, + GroupVersion: fmt.Sprintf("%s/%s", group, version), + APIResources: []metav1.APIResource{}, + } + d.resources[gvKey] = resourceList + } + + // Add the resource to the list + apiResource := metav1.APIResource{ + Name: crd.Spec.Names.Plural, + SingularName: crd.Spec.Names.Singular, + Namespaced: crd.Spec.Scope == apiextensionsv1.NamespaceScoped, + Kind: crd.Spec.Names.Kind, + // TODO(@konsalex): Can this be dynamic ever? + // Need to validate + Verbs: []string{"create", "delete", "deletecollection", "get", "list", "patch", "update", "watch"}, + ShortNames: crd.Spec.Names.ShortNames, + } + + // Add status subresource if defined + if crd.Spec.Versions[0].Subresources != nil && crd.Spec.Versions[0].Subresources.Status != nil { + apiResource.Verbs = append(apiResource.Verbs, "update") + } + + // Check if resource already exists + found := false + for i, r := range resourceList.APIResources { + if r.Name == apiResource.Name { + resourceList.APIResources[i] = apiResource + found = true + break + } + } + if !found { + resourceList.APIResources = append(resourceList.APIResources, apiResource) + } +} + +// GetAPIGroupList returns the list of API groups for /apis discovery +func (d *DiscoveryManager) GetAPIGroupList() *metav1.APIGroupList { + d.mu.RLock() + defer d.mu.RUnlock() + + groups := make([]metav1.APIGroup, 0, len(d.apiGroups)) + for _, group := range d.apiGroups { + groups = append(groups, *group) + } + + return &metav1.APIGroupList{ + TypeMeta: metav1.TypeMeta{ + Kind: "APIGroupList", + APIVersion: "v1", + }, + Groups: groups, + } +} + +// GetAPIGroup returns a specific API group +func (d *DiscoveryManager) GetAPIGroup(name string) *metav1.APIGroup { + d.mu.RLock() + defer d.mu.RUnlock() + + return d.apiGroups[name] +} + +// GetAPIResourceList returns the resource list for a specific group/version +func (d *DiscoveryManager) GetAPIResourceList(gv schema.GroupVersion) *metav1.APIResourceList { + d.mu.RLock() + defer d.mu.RUnlock() + + key := fmt.Sprintf("%s/%s", gv.Group, gv.Version) + return d.resources[key] +} + +// ServeAPIGroup handles requests to /apis/ +func (d *DiscoveryManager) ServeAPIGroup(w http.ResponseWriter, req *http.Request, group string) { + apiGroup := d.GetAPIGroup(group) + if apiGroup == nil { + http.NotFound(w, req) + return + } + + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(apiGroup) +} + +// ServeAPIResourceList handles requests to /apis// +func (d *DiscoveryManager) ServeAPIResourceList(w http.ResponseWriter, req *http.Request, gv schema.GroupVersion) { + resourceList := d.GetAPIResourceList(gv) + if resourceList == nil { + http.NotFound(w, req) + return + } + + w.Header().Set("Content-Type", "application/json") + json.NewEncoder(w).Encode(resourceList) +} diff --git a/pkg/registry/apis/apiextensions/dynamic_handler.go b/pkg/registry/apis/apiextensions/dynamic_handler.go new file mode 100644 index 00000000000..bb237bd0497 --- /dev/null +++ b/pkg/registry/apis/apiextensions/dynamic_handler.go @@ -0,0 +1,451 @@ +package apiextensions + +import ( + "context" + "fmt" + "io" + "mime" + "net/http" + "strings" + "sync" + + apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/serializer" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apiserver/pkg/endpoints/handlers/negotiation" + "k8s.io/apiserver/pkg/endpoints/request" + "k8s.io/apiserver/pkg/registry/rest" +) + +// DynamicCRHandler handles HTTP requests for dynamically registered custom resources +type DynamicCRHandler struct { + mu sync.RWMutex + // Map of group -> version -> resource -> storage + storageMap map[string]map[string]map[string]rest.Storage + // Map of group -> version -> resource -> scope (Namespaced or Cluster) + scopeMap map[string]map[string]map[string]apiextensionsv1.ResourceScope + scheme *runtime.Scheme + codecs serializer.CodecFactory +} + +// NewDynamicCRHandler creates a new dynamic custom resource handler +func NewDynamicCRHandler(scheme *runtime.Scheme) *DynamicCRHandler { + return &DynamicCRHandler{ + storageMap: make(map[string]map[string]map[string]rest.Storage), + scheme: scheme, + codecs: serializer.NewCodecFactory(scheme), + } +} + +type registeredCR struct { + storage rest.Storage + scope apiextensionsv1.ResourceScope +} + +// RegisterCustomResource registers a custom resource storage for dynamic routing +func (h *DynamicCRHandler) RegisterCustomResource( + crd *apiextensionsv1.CustomResourceDefinition, + version string, + storage rest.Storage, +) { + h.mu.Lock() + defer h.mu.Unlock() + + group := crd.Spec.Group + resource := crd.Spec.Names.Plural + + if h.storageMap[group] == nil { + h.storageMap[group] = make(map[string]map[string]rest.Storage) + } + if h.storageMap[group][version] == nil { + h.storageMap[group][version] = make(map[string]rest.Storage) + } + + h.storageMap[group][version][resource] = storage + + // Store the scope information + if h.scopeMap == nil { + h.scopeMap = make(map[string]map[string]map[string]apiextensionsv1.ResourceScope) + } + if h.scopeMap[group] == nil { + h.scopeMap[group] = make(map[string]map[string]apiextensionsv1.ResourceScope) + } + if h.scopeMap[group][version] == nil { + h.scopeMap[group][version] = make(map[string]apiextensionsv1.ResourceScope) + } + h.scopeMap[group][version][resource] = crd.Spec.Scope + + fmt.Printf("DynamicCRHandler: Registered %s/%s/%s (scope: %s)\n", group, version, resource, crd.Spec.Scope) +} + +// ServeHTTP handles HTTP requests for custom resources +func (h *DynamicCRHandler) ServeHTTP(w http.ResponseWriter, req *http.Request) { + // We could use a custom regex pattern, but RequestInfo is + // already available and gives us the match info for free. + requestInfo, ok := request.RequestInfoFrom(req.Context()) + if !ok || requestInfo == nil { + http.NotFound(w, req) + return + } + + group := requestInfo.APIGroup + version := requestInfo.APIVersion + namespace := requestInfo.Namespace + resource := requestInfo.Resource + name := requestInfo.Name + + fmt.Printf("DynamicCRHandler: %s %s (group=%s, version=%s, resource=%s, namespace=%s, name=%s)\n", + req.Method, req.URL.Path, group, version, resource, namespace, name) + + // Look up the storage + h.mu.RLock() + versionMap, groupExists := h.storageMap[group] + if !groupExists { + h.mu.RUnlock() + http.NotFound(w, req) + return + } + + resourceMap, versionExists := versionMap[version] + if !versionExists { + h.mu.RUnlock() + http.NotFound(w, req) + return + } + + storage, resourceExists := resourceMap[resource] + if !resourceExists { + h.mu.RUnlock() + http.NotFound(w, req) + return + } + + // Check scope + scope := h.scopeMap[group][version][resource] + h.mu.RUnlock() + + // Validate scope matches the request + if scope == apiextensionsv1.NamespaceScoped && namespace == "" { + http.Error(w, fmt.Sprintf("resource %s is namespace-scoped, must specify namespace", resource), http.StatusBadRequest) + return + } + if scope == apiextensionsv1.ClusterScoped && namespace != "" { + http.Error(w, fmt.Sprintf("resource %s is cluster-scoped, cannot specify namespace", resource), http.StatusBadRequest) + return + } + + // Add request info to context + ctx := req.Context() + ctx = request.WithNamespace(ctx, namespace) + ctx = request.WithRequestInfo(ctx, &request.RequestInfo{ + IsResourceRequest: true, + Path: req.URL.Path, + Verb: strings.ToLower(req.Method), + APIGroup: group, + APIVersion: version, + Namespace: namespace, + Resource: resource, + Name: name, + }) + + // Handle the request based on the storage interface + h.handleStorageRequest(ctx, w, req, storage, name, namespace) +} + +// handleStorageRequest dispatches to the appropriate storage method +func (h *DynamicCRHandler) handleStorageRequest( + ctx context.Context, + w http.ResponseWriter, + req *http.Request, + storage rest.Storage, + name string, + namespace string, +) { + // Determine media type for response + _, serializerInfo, err := negotiation.NegotiateOutputMediaType(req, h.codecs, negotiation.DefaultEndpointRestrictions) + if err != nil { + http.Error(w, fmt.Sprintf("failed to negotiate media type: %v", err), http.StatusNotAcceptable) + return + } + serializer := serializerInfo.Serializer + + switch req.Method { + case http.MethodGet: + if name == "" { + // List operation + if lister, ok := storage.(rest.Lister); ok { + h.handleList(ctx, w, req, lister, serializer) + } else { + http.Error(w, "list not supported", http.StatusMethodNotAllowed) + } + } else { + // Get operation + if getter, ok := storage.(rest.Getter); ok { + h.handleGet(ctx, w, req, getter, name, serializer) + } else { + http.Error(w, "get not supported", http.StatusMethodNotAllowed) + } + } + + case http.MethodPost: + // Create operation + if creater, ok := storage.(rest.Creater); ok { + h.handleCreate(ctx, w, req, creater, serializer) + } else { + http.Error(w, "create not supported", http.StatusMethodNotAllowed) + } + + case http.MethodPut: + // Update operation + if updater, ok := storage.(rest.Updater); ok { + h.handleUpdate(ctx, w, req, updater, name, serializer) + } else { + http.Error(w, "update not supported", http.StatusMethodNotAllowed) + } + + case http.MethodPatch: + // Patch operation + if patcher, ok := storage.(rest.Patcher); ok { + h.handlePatch(ctx, w, req, patcher, name, serializer) + } else { + http.Error(w, "patch not supported", http.StatusMethodNotAllowed) + } + + case http.MethodDelete: + // Delete operation + if deleter, ok := storage.(rest.GracefulDeleter); ok { + h.handleDelete(ctx, w, req, deleter, name, serializer) + } else { + http.Error(w, "delete not supported", http.StatusMethodNotAllowed) + } + + default: + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + } +} + +// Helper methods for each operation +func (h *DynamicCRHandler) handleList(ctx context.Context, w http.ResponseWriter, req *http.Request, lister rest.Lister, serializer runtime.Serializer) { + // TODO(@konsalex): Parse list options from query parameters + obj, err := lister.List(ctx, nil) + if err != nil { + http.Error(w, fmt.Sprintf("failed to list: %v", err), http.StatusInternalServerError) + return + } + + h.writeResponse(w, obj, serializer) +} + +func (h *DynamicCRHandler) handleGet(ctx context.Context, w http.ResponseWriter, req *http.Request, getter rest.Getter, name string, serializer runtime.Serializer) { + // TODO(@konsalex): Parse get options + obj, err := getter.Get(ctx, name, nil) + if err != nil { + http.Error(w, fmt.Sprintf("failed to get: %v", err), http.StatusInternalServerError) + return + } + + h.writeResponse(w, obj, serializer) +} + +func (h *DynamicCRHandler) handleCreate(ctx context.Context, w http.ResponseWriter, req *http.Request, creater rest.Creater, serializer runtime.Serializer) { + // Read the request body + body, err := io.ReadAll(req.Body) + if err != nil { + http.Error(w, fmt.Sprintf("failed to read body: %v", err), http.StatusBadRequest) + return + } + defer req.Body.Close() + + fmt.Printf("DynamicCRHandler: handleCreate received body: %s\n", string(body)) + + // Decode the body + decoder := h.codecs.UniversalDeserializer() + obj, gvk, err := decoder.Decode(body, nil, nil) + if err != nil { + http.Error(w, fmt.Sprintf("failed to decode body: %v", err), http.StatusBadRequest) + return + } + + fmt.Printf("DynamicCRHandler: Decoded object type: %T, GVK: %v\n", obj, gvk) + + // Ensure we have an unstructured object + if obj == nil { + http.Error(w, "decoded object is nil", http.StatusBadRequest) + return + } + + // Create the object + fmt.Printf("DynamicCRHandler: Calling creater.Create...\n") + created, err := creater.Create(ctx, obj, nil, nil) + if err != nil { + fmt.Printf("DynamicCRHandler: Create failed: %v\n", err) + http.Error(w, fmt.Sprintf("failed to create: %v", err), http.StatusInternalServerError) + return + } + + fmt.Printf("DynamicCRHandler: Create succeeded, writing response\n") + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusCreated) + if err := serializer.Encode(created, w); err != nil { + fmt.Printf("DynamicCRHandler: Failed to encode response: %v\n", err) + } +} + +func (h *DynamicCRHandler) handleUpdate(ctx context.Context, w http.ResponseWriter, req *http.Request, updater rest.Updater, name string, serializer runtime.Serializer) { + fmt.Printf("DynamicCRHandler: handleUpdate for resource: %s\n", name) + + // Read the request body + body, err := io.ReadAll(req.Body) + if err != nil { + http.Error(w, fmt.Sprintf("failed to read body: %v", err), http.StatusBadRequest) + return + } + defer req.Body.Close() + + // Decode the body + decoder := h.codecs.UniversalDeserializer() + obj, gvk, err := decoder.Decode(body, nil, nil) + if err != nil { + http.Error(w, fmt.Sprintf("failed to decode body: %v", err), http.StatusBadRequest) + return + } + + fmt.Printf("DynamicCRHandler: Decoded object type: %T, GVK: %v\n", obj, gvk) + + // Create UpdatedObjectInfo + objInfo := rest.DefaultUpdatedObjectInfo(obj) + + // Update the object + updated, created, err := updater.Update(ctx, name, objInfo, nil, nil, false, nil) + if err != nil { + fmt.Printf("DynamicCRHandler: Update failed: %v\n", err) + if apierrors.IsNotFound(err) { + http.Error(w, fmt.Sprintf("resource not found: %v", err), http.StatusNotFound) + } else { + http.Error(w, fmt.Sprintf("failed to update: %v", err), http.StatusInternalServerError) + } + return + } + + w.Header().Set("Content-Type", "application/json") + if created { + w.WriteHeader(http.StatusCreated) + } else { + w.WriteHeader(http.StatusOK) + } + + fmt.Printf("DynamicCRHandler: Update succeeded\n") + if err := serializer.Encode(updated, w); err != nil { + fmt.Printf("DynamicCRHandler: Failed to encode response: %v\n", err) + } +} + +func (h *DynamicCRHandler) handlePatch(ctx context.Context, w http.ResponseWriter, req *http.Request, patcher rest.Patcher, name string, serializer runtime.Serializer) { + // rest.Patcher is actually just Getter + Updater + // We need to get the existing object, apply the patch, then update + + // Read the request body + patchBytes, err := io.ReadAll(req.Body) + if err != nil { + http.Error(w, fmt.Sprintf("failed to read body: %v", err), http.StatusBadRequest) + return + } + defer req.Body.Close() + + // Get the existing object to verify it exists + _, err = patcher.Get(ctx, name, nil) + if err != nil { + if apierrors.IsNotFound(err) { + http.Error(w, fmt.Sprintf("resource not found: %v", err), http.StatusNotFound) + } else { + http.Error(w, fmt.Sprintf("failed to get resource: %v", err), http.StatusInternalServerError) + } + return + } + + // Determine patch type from Content-Type header + // Using mime here as content type might be like: + // "Content-Type: application/json; charset=utf-8" + contentType, _, err := mime.ParseMediaType(req.Header.Get("Content-Type")) + if err != nil { + http.Error(w, fmt.Sprintf("error parsing Content-Type: %s", contentType), http.StatusUnsupportedMediaType) + return + } + + var patchType types.PatchType + + switch contentType { + case string(types.JSONPatchType): + patchType = types.JSONPatchType + case string(types.MergePatchType): + patchType = types.MergePatchType + case string(types.StrategicMergePatchType): + patchType = types.StrategicMergePatchType + case "application/json": + // Default to merge patch for plain JSON + // We cannot fall back to Strategic as we miss + // Go structs for dynamic CRDs + patchType = types.MergePatchType + default: + http.Error(w, fmt.Sprintf("unsupported Content-Type: %s", contentType), http.StatusUnsupportedMediaType) + return + } + + // Apply the patch via our custom storage (which has the Patch method) + if storage, ok := patcher.(*customResourceStorage); ok { + patched, err := storage.Patch(ctx, name, patchType, patchBytes, nil) + if err != nil { + if apierrors.IsNotFound(err) { + http.Error(w, fmt.Sprintf("resource not found: %v", err), http.StatusNotFound) + } else { + http.Error(w, fmt.Sprintf("failed to patch: %v", err), http.StatusInternalServerError) + } + return + } + + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + if err := serializer.Encode(patched, w); err != nil { + fmt.Printf("DynamicCRHandler: Failed to encode response: %v\n", err) + } + } else { + http.Error(w, "patch not supported for this resource type", http.StatusNotImplemented) + } +} + +func (h *DynamicCRHandler) handleDelete(ctx context.Context, w http.ResponseWriter, req *http.Request, deleter rest.GracefulDeleter, name string, serializer runtime.Serializer) { + // Parse delete options from query/body + deleteOptions := &metav1.DeleteOptions{} + + // Delete the object + obj, deleted, err := deleter.Delete(ctx, name, nil, deleteOptions) + if err != nil { + if apierrors.IsNotFound(err) { + http.Error(w, fmt.Sprintf("resource not found: %v", err), http.StatusNotFound) + } else { + http.Error(w, fmt.Sprintf("failed to delete: %v", err), http.StatusInternalServerError) + } + return + } + + w.Header().Set("Content-Type", "application/json") + if deleted { + w.WriteHeader(http.StatusOK) + } else { + // Return 202 Accepted if deletion is pending (e.g., finalizers) + w.WriteHeader(http.StatusAccepted) + } + if err := serializer.Encode(obj, w); err != nil { + fmt.Printf("DynamicCRHandler: Failed to encode response: %v\n", err) + } +} + +func (h *DynamicCRHandler) writeResponse(w http.ResponseWriter, obj runtime.Object, serializer runtime.Serializer) { + w.Header().Set("Content-Type", "application/json") + if err := serializer.Encode(obj, w); err != nil { + http.Error(w, fmt.Sprintf("failed to encode response: %v", err), http.StatusInternalServerError) + } +} diff --git a/pkg/registry/apis/apiextensions/dynamic_registry.go b/pkg/registry/apis/apiextensions/dynamic_registry.go new file mode 100644 index 00000000000..e01553b1682 --- /dev/null +++ b/pkg/registry/apis/apiextensions/dynamic_registry.go @@ -0,0 +1,452 @@ +package apiextensions + +import ( + "context" + "fmt" + "net/http" + "sync" + + authlib "github.com/grafana/authlib/types" + apidiscoveryv2 "k8s.io/api/apidiscovery/v2" + apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + "k8s.io/apiextensions-apiserver/pkg/controller/openapi/builder" + apiextensionsfeatures "k8s.io/apiextensions-apiserver/pkg/features" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apiserver/pkg/endpoints/discovery/aggregated" + "k8s.io/apiserver/pkg/registry/generic" + genericregistry "k8s.io/apiserver/pkg/registry/generic/registry" + "k8s.io/apiserver/pkg/registry/rest" + genericapiserver "k8s.io/apiserver/pkg/server" + utilfeature "k8s.io/apiserver/pkg/util/feature" + "k8s.io/kube-openapi/pkg/spec3" + + "github.com/grafana/grafana/pkg/storage/unified/resource" +) + +// DynamicRegistry manages the dynamic registration and unregistration of custom resources +type DynamicRegistry struct { + scheme *runtime.Scheme + optsGetter generic.RESTOptionsGetter + unifiedClient resource.ResourceClient + apiGroupInfo *genericapiserver.APIGroupInfo + accessClient authlib.AccessClient + server *genericapiserver.GenericAPIServer // API server to install new groups + discoveryManager *DiscoveryManager + + mu sync.RWMutex + registrations map[string]*customResourceRegistration // key: group/version/resource + apiGroups map[string]*genericapiserver.APIGroupInfo // key: group name + + // openAPISpecs caches OpenAPI specs per GroupVersion and CRD Name + openAPISpecs map[schema.GroupVersion]map[string]*spec3.OpenAPI +} + +type customResourceRegistration struct { + crd *apiextensionsv1.CustomResourceDefinition + storage *customResourceStorage +} + +// NewDynamicRegistry creates a new dynamic registry +func NewDynamicRegistry( + scheme *runtime.Scheme, + optsGetter generic.RESTOptionsGetter, + apiGroupInfo *genericapiserver.APIGroupInfo, + accessClient authlib.AccessClient, +) *DynamicRegistry { + return &DynamicRegistry{ + scheme: scheme, + optsGetter: optsGetter, + apiGroupInfo: apiGroupInfo, + accessClient: accessClient, + discoveryManager: NewDiscoveryManager(), + registrations: make(map[string]*customResourceRegistration), + apiGroups: make(map[string]*genericapiserver.APIGroupInfo), + openAPISpecs: make(map[schema.GroupVersion]map[string]*spec3.OpenAPI), + } +} + +// SetAPIServer sets the API server for dynamic group installation +func (r *DynamicRegistry) SetAPIServer(server *genericapiserver.GenericAPIServer) { + r.mu.Lock() + defer r.mu.Unlock() + r.server = server + + // Register discovery handlers for any API groups that were registered before we had the server reference + for groupName := range r.apiGroups { + // Find the version for this group + for _, reg := range r.registrations { + if reg.crd.Spec.Group == groupName { + versionName := reg.crd.Spec.Versions[0].Name + r.registerDiscoveryHandlers(groupName, versionName) + break + } + } + } + + // Wrap the main /apis handler to include our custom groups + r.registerAllWithAggregatedDiscovery() + + // Note: OpenAPI registration is deferred to a PostStartHook because + // OpenAPIV3VersionedService is only available after PrepareRun() +} + +// RegisterOpenAPIForExistingCRDs registers OpenAPI specs for all existing CRDs +// This is called from a PostStartHook after the server is fully prepared +func (r *DynamicRegistry) RegisterOpenAPIForExistingCRDs() error { + r.mu.Lock() + defer r.mu.Unlock() + + // Register OpenAPI specs for all CRDs that were registered before the server was fully prepared + for _, reg := range r.registrations { + if err := r.updateOpenAPISpecLocked(reg.crd); err != nil { + fmt.Printf("Warning: failed to update OpenAPI spec for CRD %s: %v\n", reg.crd.Name, err) + } + } + + return nil +} + +// registerDiscoveryHandlers registers HTTP handlers for discovery endpoints +func (r *DynamicRegistry) registerDiscoveryHandlers(group, version string) { + if r.server == nil { + return + } + + // Register /apis/ handler + groupPath := fmt.Sprintf("/apis/%s", group) + r.server.Handler.NonGoRestfulMux.HandleFunc(groupPath, func(w http.ResponseWriter, req *http.Request) { + r.discoveryManager.ServeAPIGroup(w, req, group) + }) + // DEBUG + fmt.Printf("Registered discovery handler: %s\n", groupPath) + + // Register /apis// handler + gv := schema.GroupVersion{Group: group, Version: version} + versionPath := fmt.Sprintf("/apis/%s/%s", group, version) + r.server.Handler.NonGoRestfulMux.HandleFunc(versionPath, func(w http.ResponseWriter, req *http.Request) { + r.discoveryManager.ServeAPIResourceList(w, req, gv) + }) + fmt.Printf("Registered discovery handler: %s\n", versionPath) +} + +// registerAllWithAggregatedDiscovery registers all custom groups with the server's aggregated discovery manager +func (r *DynamicRegistry) registerAllWithAggregatedDiscovery() { + if r.server == nil || r.server.AggregatedDiscoveryGroupManager == nil { + return + } + + // Register each custom group with the discovery manager + for _, apiGroup := range r.discoveryManager.apiGroups { + for _, gvDiscovery := range apiGroup.Versions { + r.registerWithAggregatedDiscovery(apiGroup.Name, gvDiscovery.Version) + } + } +} + +// registerWithAggregatedDiscovery registers a specific group version with the aggregated discovery manager +func (r *DynamicRegistry) registerWithAggregatedDiscovery(groupName, versionName string) { + if r.server == nil || r.server.AggregatedDiscoveryGroupManager == nil { + return + } + + // Get the resources for this version + gvKey := fmt.Sprintf("%s/%s", groupName, versionName) + resourceList := r.discoveryManager.resources[gvKey] + + if resourceList != nil { + // Convert our metav1.APIResourceList to apidiscoveryv2.APIVersionDiscovery + apiResources := make([]apidiscoveryv2.APIResourceDiscovery, 0, len(resourceList.APIResources)) + for _, res := range resourceList.APIResources { + apiRes := apidiscoveryv2.APIResourceDiscovery{ + Resource: res.Name, + ResponseKind: &metav1.GroupVersionKind{ + Group: groupName, + Version: versionName, + Kind: res.Kind, + }, + Scope: getScopeType(res.Namespaced), + ShortNames: res.ShortNames, + Verbs: res.Verbs, + } + apiResources = append(apiResources, apiRes) + } + + versionDiscovery := apidiscoveryv2.APIVersionDiscovery{ + Version: versionName, + Resources: apiResources, + } + + // Create a manager with CRD source + manager := r.server.AggregatedDiscoveryGroupManager.WithSource(aggregated.CRDSource) + + // Add group version + manager.AddGroupVersion( + groupName, + versionDiscovery, + ) + + // Set priority to ensure it's discoverable + gv := metav1.GroupVersion{ + Group: groupName, + Version: versionName, + } + manager.SetGroupVersionPriority(gv, 1000, 100) + + fmt.Printf("✅ Registered %s with AggregatedDiscoveryGroupManager (Source: CRD)\n", gvKey) + } +} + +// getScopeType converts boolean namespaced flag to apidiscoveryv2 scope type +func getScopeType(namespaced bool) apidiscoveryv2.ResourceScope { + if namespaced { + return apidiscoveryv2.ScopeNamespace + } + return apidiscoveryv2.ScopeCluster +} + +// Start begins watching CRDs for dynamic registration +func (r *DynamicRegistry) Start(ctx context.Context, crdStore *genericregistry.Store) { + // Currently we don't watch for changes + // CRDs are loaded during UpdateAPIGroupInfo before the server starts + // We need to implement proper watching mechanism here. + // TODO(@konsalex): Watch for CRD changes and call RegisterCRD/UpdateCRD/UnregisterCRD + <-ctx.Done() +} + +// RegisterCRD registers a custom resource dynamically based on the CRD spec +func (r *DynamicRegistry) RegisterCRD(crd *apiextensionsv1.CustomResourceDefinition) error { + r.mu.Lock() + defer r.mu.Unlock() + + // TODO(@konsalex): Support multiple versions + if len(crd.Spec.Versions) != 1 { + return fmt.Errorf("only single-version CRDs are supported") + } + + version := crd.Spec.Versions[0] + + // Create storage for the custom resource + crStorage, err := NewCustomResourceStorage( + crd, + version.Name, + r.scheme, + r.optsGetter, + r.accessClient, + r.unifiedClient, + ) + if err != nil { + return fmt.Errorf("failed to create custom resource storage: %w", err) + } + + // Get or create API group info for this custom resource's group + group := crd.Spec.Group + apiGroupInfo, ok := r.apiGroups[group] + isNewGroup := !ok + if !ok { + // Create a simple API group info without using NewDefaultAPIGroupInfo + // to avoid OpenAPI validation issues with unstructured types + gvInfo := genericapiserver.APIGroupInfo{ + PrioritizedVersions: []schema.GroupVersion{{Group: group, Version: version.Name}}, + VersionedResourcesStorageMap: make(map[string]map[string]rest.Storage), + Scheme: r.scheme, + NegotiatedSerializer: r.apiGroupInfo.NegotiatedSerializer, + ParameterCodec: r.apiGroupInfo.ParameterCodec, + } + + r.apiGroups[group] = &gvInfo + apiGroupInfo = &gvInfo + } + + // Get or create the storage map for this version + storageMap, ok := apiGroupInfo.VersionedResourcesStorageMap[version.Name] + if !ok { + storageMap = make(map[string]rest.Storage) + apiGroupInfo.VersionedResourcesStorageMap[version.Name] = storageMap + + // Update prioritized versions if this is a new version + found := false + for _, gv := range apiGroupInfo.PrioritizedVersions { + if gv.Version == version.Name { + found = true + break + } + } + if !found { + apiGroupInfo.PrioritizedVersions = append(apiGroupInfo.PrioritizedVersions, + schema.GroupVersion{Group: group, Version: version.Name}) + } + } + + // Register the custom resource storage + resourcePath := crd.Spec.Names.Plural + storageMap[resourcePath] = crStorage + + // Register status subresource if defined + // TODO(@konsalex): Implement status subresource for direct storage + // if version.Subresources != nil && version.Subresources.Status != nil { + // storageMap[resourcePath+"/status"] = crStorage + // } + + // Store the registration + key := fmt.Sprintf("%s/%s/%s", crd.Spec.Group, version.Name, crd.Spec.Names.Plural) + r.registrations[key] = &customResourceRegistration{ + crd: crd.DeepCopy(), + storage: crStorage, + } + + // Add to discovery manager + // (this will make kubectl work properly) + r.discoveryManager.AddCustomResource(crd) + + // Register discovery HTTP handlers if we have the server + if r.server != nil && isNewGroup { + r.registerDiscoveryHandlers(group, version.Name) + // Also register with aggregated discovery manager + r.registerWithAggregatedDiscovery(group, version.Name) + } + + // Update OpenAPI spec + if err := r.updateOpenAPISpecLocked(crd); err != nil { + fmt.Printf("Warning: failed to update OpenAPI spec for CRD %s: %v\n", crd.Name, err) + } + + return nil +} + +// UpdateCRD updates a custom resource registration +func (r *DynamicRegistry) UpdateCRD(crd *apiextensionsv1.CustomResourceDefinition) error { + r.mu.Lock() + defer r.mu.Unlock() + + // Hack for now, we just un-register it and the re-register it. + // Not sure if there is any drawback + if err := r.unregisterCRDLocked(crd); err != nil { + return err + } + + return r.RegisterCRD(crd) +} + +// UnregisterCRD removes a custom resource registration +func (r *DynamicRegistry) UnregisterCRD(crd *apiextensionsv1.CustomResourceDefinition) error { + r.mu.Lock() + defer r.mu.Unlock() + + return r.unregisterCRDLocked(crd) +} + +func (r *DynamicRegistry) unregisterCRDLocked(crd *apiextensionsv1.CustomResourceDefinition) error { + if len(crd.Spec.Versions) != 1 { + return fmt.Errorf("only single-version CRDs are supported") + } + + version := crd.Spec.Versions[0] + key := fmt.Sprintf("%s/%s/%s", crd.Spec.Group, version.Name, crd.Spec.Names.Plural) + + // Remove from registrations + delete(r.registrations, key) + + // Remove from API group info storage map + if storageMap, ok := r.apiGroupInfo.VersionedResourcesStorageMap[version.Name]; ok { + resourcePath := crd.Spec.Names.Plural + delete(storageMap, resourcePath) + delete(storageMap, resourcePath+"/status") + } + + // Remove from OpenAPI specs + gv := schema.GroupVersion{Group: crd.Spec.Group, Version: version.Name} + if specs, ok := r.openAPISpecs[gv]; ok { + delete(specs, crd.Name) + if len(specs) == 0 { + delete(r.openAPISpecs, gv) + // Remove from service + if r.server != nil && r.server.OpenAPIV3VersionedService != nil { + path := fmt.Sprintf("apis/%s/%s", gv.Group, gv.Version) + r.server.OpenAPIV3VersionedService.DeleteGroupVersion(path) + } + } else { + // Update with remaining specs + if err := r.updateGroupVersionOpenAPILocked(gv); err != nil { + return fmt.Errorf("failed to update OpenAPI spec after unregistering CRD: %w", err) + } + } + } + + return nil +} + +// GetRegistration returns the registration for a custom resource +func (r *DynamicRegistry) GetRegistration(group, version, resource string) (*customResourceRegistration, bool) { + r.mu.RLock() + defer r.mu.RUnlock() + + key := fmt.Sprintf("%s/%s/%s", group, version, resource) + reg, ok := r.registrations[key] + return reg, ok +} + +// updateOpenAPISpecLocked builds and updates the OpenAPI spec for the given CRD +// mu must be held by caller +func (r *DynamicRegistry) updateOpenAPISpecLocked(crd *apiextensionsv1.CustomResourceDefinition) error { + if r.server == nil || r.server.OpenAPIV3VersionedService == nil { + return nil + } + + for _, v := range crd.Spec.Versions { + if !v.Served { + continue + } + + // Build OpenAPI V3 spec + spec, err := builder.BuildOpenAPIV3(crd, v.Name, builder.Options{ + V2: false, + IncludeSelectableFields: utilfeature.DefaultFeatureGate.Enabled(apiextensionsfeatures.CustomResourceFieldSelectors), + }) + if err != nil { + return fmt.Errorf("failed to build OpenAPI V3 spec: %w", err) + } + + gv := schema.GroupVersion{Group: crd.Spec.Group, Version: v.Name} + if r.openAPISpecs[gv] == nil { + r.openAPISpecs[gv] = make(map[string]*spec3.OpenAPI) + } + r.openAPISpecs[gv][crd.Name] = spec + + // Update the group version spec in the service + if err := r.updateGroupVersionOpenAPILocked(gv); err != nil { + return err + } + } + + return nil +} + +// updateGroupVersionOpenAPILocked merges all specs for a GV and updates the service +// mu must be held by caller +func (r *DynamicRegistry) updateGroupVersionOpenAPILocked(gv schema.GroupVersion) error { + if r.server == nil || r.server.OpenAPIV3VersionedService == nil { + return nil + } + + specsMap := r.openAPISpecs[gv] + if len(specsMap) == 0 { + return nil + } + + var specs []*spec3.OpenAPI + for _, spec := range specsMap { + specs = append(specs, spec) + } + + mergedSpec, err := builder.MergeSpecsV3(specs...) + if err != nil { + return fmt.Errorf("failed to merge specs: %w", err) + } + + path := fmt.Sprintf("apis/%s/%s", gv.Group, gv.Version) + r.server.OpenAPIV3VersionedService.UpdateGroupVersion(path, mergedSpec) + + return nil +} diff --git a/pkg/registry/apis/apiextensions/register.go b/pkg/registry/apis/apiextensions/register.go new file mode 100644 index 00000000000..e6e95eea193 --- /dev/null +++ b/pkg/registry/apis/apiextensions/register.go @@ -0,0 +1,362 @@ +package apiextensions + +import ( + "context" + "fmt" + "net/http" + + "github.com/prometheus/client_golang/prometheus" + apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + apiextensionsopenapi "k8s.io/apiextensions-apiserver/pkg/generated/openapi" + metainternalversion "k8s.io/apimachinery/pkg/apis/meta/internalversion" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + genericregistry "k8s.io/apiserver/pkg/registry/generic/registry" + "k8s.io/apiserver/pkg/registry/rest" + genericapiserver "k8s.io/apiserver/pkg/server" + "k8s.io/kube-openapi/pkg/common" + + authlib "github.com/grafana/authlib/types" + + "github.com/grafana/grafana/pkg/apimachinery/identity" + "github.com/grafana/grafana/pkg/apimachinery/utils" + grafanaregistry "github.com/grafana/grafana/pkg/apiserver/registry/generic" + "github.com/grafana/grafana/pkg/services/apiserver/builder" + "github.com/grafana/grafana/pkg/services/featuremgmt" + "github.com/grafana/grafana/pkg/setting" + "github.com/grafana/grafana/pkg/storage/unified/apistore" + "github.com/grafana/grafana/pkg/storage/unified/resource" +) + +var _ builder.APIGroupBuilder = (*APIExtensionsBuilder)(nil) + +// APIExtensionsBuilder implements builder.APIGroupBuilder for CustomResourceDefinitions +type APIExtensionsBuilder struct { + features featuremgmt.FeatureToggles + storage *genericregistry.Store + accessClient authlib.AccessClient + dynamicReg *DynamicRegistry + restOptGetter *apistore.RESTOptionsGetter + apiregistrar builder.APIRegistrar + unifiedClient resource.ResourceClient + preloadedCRDs []*apiextensionsv1.CustomResourceDefinition // CRDs loaded before init + dynamicHandler *DynamicCRHandler // Dynamic handler for custom resources + server *genericapiserver.GenericAPIServer // The running API server +} + +// SetAPIServer sets the API server instance +// This allows the builder to register dynamic API groups (CRDs) +func (b *APIExtensionsBuilder) SetAPIServer(server *genericapiserver.GenericAPIServer) { + b.server = server + if b.dynamicReg != nil { + b.dynamicReg.SetAPIServer(server) + } + + // Register a PostStartHook to register OpenAPI specs after the server is fully prepared + // This is necessary because OpenAPIV3VersionedService is only available after PrepareRun() + server.AddPostStartHookOrDie("apiextensions-openapi", func(context genericapiserver.PostStartHookContext) error { + if b.dynamicReg != nil { + return b.dynamicReg.RegisterOpenAPIForExistingCRDs() + } + return nil + }) +} + +// RegisterAPIService registers the apiextensions API group in single-tenant mode +func RegisterAPIService( + cfg *setting.Cfg, + features featuremgmt.FeatureToggles, + apiregistration builder.APIRegistrar, + accessClient authlib.AccessClient, + registerer prometheus.Registerer, + unified resource.ResourceClient, +) (*APIExtensionsBuilder, error) { + if !features.IsEnabledGlobally(featuremgmt.FlagApiExtensions) { + return nil, fmt.Errorf("apiextensions feature flag is not enabled") + } + + b := &APIExtensionsBuilder{ + features: features, + accessClient: accessClient, + apiregistrar: apiregistration, + unifiedClient: unified, + } + + // Register the CRD API group + apiregistration.RegisterAPI(b) + + // NOTE: We can't load CRDs here because unified storage isn't fully initialized yet + // We'll use a PostStartHook instead to load CRDs after the server is running + // See the postStartHook field and RegisterPostStartHooks method below + + return b, nil +} + +// loadAndRegisterCRDsWithDynamicHandler loads CRDs from storage and registers their handlers +// This is called during UpdateAPIGroupInfo when storage IS ready +func (b *APIExtensionsBuilder) loadAndRegisterCRDsWithDynamicHandler( + ctx context.Context, + crdStore *genericregistry.Store, + opts builder.APIGroupOptions, +) error { + // Create a system context with fallback auth + systemCtx := resource.WithFallback(ctx) + systemCtx = authlib.WithAuthInfo(systemCtx, &identity.StaticRequester{ + Type: authlib.TypeServiceAccount, + Login: "system:apiextensions", + UserID: 0, + UserUID: "system:apiextensions", + OrgID: 1, + OrgRole: identity.RoleAdmin, + IsGrafanaAdmin: true, + }) + + // List all CRDs using the initialized storage + listObj, err := crdStore.List(systemCtx, &metainternalversion.ListOptions{}) + if err != nil { + return fmt.Errorf("failed to list CRDs: %w", err) + } + + crdList, ok := listObj.(*apiextensionsv1.CustomResourceDefinitionList) + if !ok { + return fmt.Errorf("unexpected list type: %T", listObj) + } + + if len(crdList.Items) == 0 { + fmt.Println("No existing CRDs found in storage") + return nil + } + + fmt.Printf("Found %d CRDs in storage, registering with dynamic handler...\n", len(crdList.Items)) + + // For each CRD, create storage and register it with the dynamic handler + for i := range crdList.Items { + crd := &crdList.Items[i] + + fmt.Printf(" - Processing CRD: %s (group: %s, version: %s, resource: %s)\n", + crd.Name, crd.Spec.Group, crd.Spec.Versions[0].Name, crd.Spec.Names.Plural) + + // Support only single version for now + if len(crd.Spec.Versions) != 1 { + fmt.Printf(" Warning: only single-version CRDs supported, skipping %s\n", crd.Name) + continue + } + + version := crd.Spec.Versions[0] + + // Create storage for this custom resource + crStorage, err := NewCustomResourceStorage( + crd, + version.Name, + opts.Scheme, + opts.OptsGetter, + b.accessClient, + b.unifiedClient, + ) + if err != nil { + fmt.Printf(" Warning: failed to create storage for CRD %s: %v\n", crd.Name, err) + continue + } + + // Register with the dynamic handler (for HTTP routing) + b.dynamicHandler.RegisterCustomResource(crd, version.Name, crStorage) + + // Register with the dynamic registry (for API discovery) + fmt.Printf(" DEBUG: Registering CRD with dynamic registry (b.dynamicReg=%v)...\n", b.dynamicReg != nil) + if err := b.dynamicReg.RegisterCRD(crd); err != nil { + fmt.Printf(" Warning: failed to register CRD with dynamic registry: %v\n", err) + } else { + fmt.Printf(" ✓ Registered CRD with dynamic registry\n") + } + + fmt.Printf(" ✓ Registered with dynamic handler: %s/%s/%s\n", + crd.Spec.Group, version.Name, crd.Spec.Names.Plural) + } + + return nil +} + +// loadAndRegisterCRDs loads existing CRDs from storage and registers their API group builders +// This is called during UpdateAPIGroupInfo when storage is ready +func (b *APIExtensionsBuilder) loadAndRegisterCRDs(ctx context.Context, crdStore *genericregistry.Store, opts builder.APIGroupOptions) error { + if b.apiregistrar == nil { + return fmt.Errorf("apiregistrar is nil") + } + + // Create a system context with fallback auth for initialization + // This allows us to list CRDs without a user session during server startup + // TODO(@konsalex): Does this cause any security issue? Not 100% how to authenticate + // a service call like this + systemCtx := resource.WithFallback(ctx) + systemCtx = authlib.WithAuthInfo(systemCtx, &identity.StaticRequester{ + Type: authlib.TypeServiceAccount, + Login: "system:apiextensions", + UserID: 0, + UserUID: "system:apiextensions", + OrgID: 1, + OrgRole: identity.RoleAdmin, + IsGrafanaAdmin: true, + }) + + // List all CRDs using the initialized storage with system context + listObj, err := crdStore.List(systemCtx, &metainternalversion.ListOptions{}) + if err != nil { + return fmt.Errorf("failed to list CRDs: %w", err) + } + + crdList, ok := listObj.(*apiextensionsv1.CustomResourceDefinitionList) + if !ok { + return fmt.Errorf("unexpected list type: %T", listObj) + } + + if len(crdList.Items) == 0 { + return nil + } + + // Register a builder for each CRD + for i := range crdList.Items { + crd := &crdList.Items[i] + + // Register with the dynamic registry for API group installation + if err := b.dynamicReg.RegisterCRD(crd); err != nil { + // TODO(@konsalex): Add logger in the context and use it + fmt.Printf(" Warning: failed to register CRD %s in dynamic registry: %v\n", crd.Name, err) + continue + } + // DEBUG + fmt.Printf("Registered custom resource API group: %s/%s (resource: %s)\n", + crd.Spec.Group, crd.Spec.Versions[0].Name, crd.Spec.Names.Plural) + } + + return nil +} + +// NewAPIService creates an APIExtensionsBuilder for multi-tenant mode +// TODO(@konsalex): NOT YET IMPLEMENTED properly +func NewAPIService( + accessClient authlib.AccessClient, + unified resource.ResourceClient, + registerer prometheus.Registerer, + features featuremgmt.FeatureToggles, +) (*APIExtensionsBuilder, error) { + return &APIExtensionsBuilder{ + features: features, + accessClient: accessClient, + }, nil +} + +func (b *APIExtensionsBuilder) GetGroupVersion() schema.GroupVersion { + return apiextensionsv1.SchemeGroupVersion +} + +func (b *APIExtensionsBuilder) InstallSchema(scheme *runtime.Scheme) error { + gv := b.GetGroupVersion() + + // We don't register CRD types with AddKnownTypes here because: + // 1. The external k8s.io/apiextensions-apiserver types don't have openapi-gen annotations + // 2. This would cause the API server to try to generate OpenAPI specs for them + // 3. We register them at the storage level in UpdateAPIGroupInfo instead + + // We only need to add the metav1 types to this group version + metav1.AddToGroupVersion(scheme, gv) + return scheme.SetVersionPriority(gv) +} + +func (b *APIExtensionsBuilder) AllowedV0Alpha1Resources() []string { + return nil +} + +// GetDynamicHandler returns the dynamic custom resource handler +// This can be used to install it as a fallback handler in the HTTP server +func (b *APIExtensionsBuilder) GetDynamicHandler() http.Handler { + if b.dynamicHandler == nil { + return nil + } + return b.dynamicHandler +} + +func (b *APIExtensionsBuilder) UpdateAPIGroupInfo( + apiGroupInfo *genericapiserver.APIGroupInfo, + opts builder.APIGroupOptions, +) error { + // Register storage options for CRDs + opts.StorageOptsRegister( + schema.GroupResource{Group: apiextensionsv1.GroupName, Resource: "customresourcedefinitions"}, + apistore.StorageOptions{}, + ) + + // Register the CRD types directly with the scheme now (at storage registration time) + // This avoids OpenAPI generation issues during schema installation + gv := b.GetGroupVersion() + opts.Scheme.AddKnownTypes(gv, + &apiextensionsv1.CustomResourceDefinition{}, + &apiextensionsv1.CustomResourceDefinitionList{}, + ) + + // Create the main CRD storage + crdResourceInfo := utils.NewResourceInfo( + apiextensionsv1.GroupName, + "v1", + "customresourcedefinitions", + "customresourcedefinition", + "CustomResourceDefinition", + func() runtime.Object { return &apiextensionsv1.CustomResourceDefinition{} }, + func() runtime.Object { return &apiextensionsv1.CustomResourceDefinitionList{} }, + utils.TableColumns{}, + ) + crdResourceInfoWithScope := crdResourceInfo.WithClusterScope() + + unified, err := grafanaregistry.NewRegistryStore( + opts.Scheme, + crdResourceInfoWithScope, + opts.OptsGetter, + ) + if err != nil { + return fmt.Errorf("failed to create CRD storage: %w", err) + } + + b.storage = unified + b.restOptGetter = opts.OptsGetter.(*apistore.RESTOptionsGetter) + + // Initialize dynamic registry for custom resources + b.dynamicReg = NewDynamicRegistry( + opts.Scheme, + opts.OptsGetter, + apiGroupInfo, + b.accessClient, + ) + if b.server != nil { + b.dynamicReg.SetAPIServer(b.server) + } + + // Create the dynamic handler for custom resources + b.dynamicHandler = NewDynamicCRHandler(opts.Scheme) + + // NOW storage is ready, load CRDs and register their handlers dynamically + if err := b.loadAndRegisterCRDsWithDynamicHandler(context.Background(), unified, opts); err != nil { + // TODO(@konsalex): use logger here + fmt.Printf("failed to load and register CRDs: %v\n", err) + // Don't fail - CRDs can be created later + } + + // Start watching CRDs for changes (WIP) + go b.dynamicReg.Start(context.Background(), unified) + + storage := map[string]rest.Storage{} + storage["customresourcedefinitions"] = &crdStorage{ + Store: unified, + dynamicReg: b.dynamicReg, + } + storage["customresourcedefinitions/status"] = &crdStatusStorage{Store: unified} + + apiGroupInfo.VersionedResourcesStorageMap[apiextensionsv1.SchemeGroupVersion.Version] = storage + + return nil +} + +func (b *APIExtensionsBuilder) GetOpenAPIDefinitions() common.GetOpenAPIDefinitions { + return func(ref common.ReferenceCallback) map[string]common.OpenAPIDefinition { + return apiextensionsopenapi.GetOpenAPIDefinitions(ref) + } +} diff --git a/pkg/registry/apis/apiextensions/resources/example-crd.yaml b/pkg/registry/apis/apiextensions/resources/example-crd.yaml new file mode 100644 index 00000000000..c44494a4954 --- /dev/null +++ b/pkg/registry/apis/apiextensions/resources/example-crd.yaml @@ -0,0 +1,69 @@ +apiVersion: apiextensions.k8s.io/v1 +kind: CustomResourceDefinition +metadata: + name: widgets.customcrdtest.grafana.app +spec: + group: customcrdtest.grafana.app + names: + kind: Widget + listKind: WidgetList + plural: widgets + singular: widget + shortNames: + - wg + scope: Namespaced + versions: + - name: v1 + served: true + storage: true + schema: + openAPIV3Schema: + type: object + properties: + apiVersion: + type: string + kind: + type: string + metadata: + type: object + spec: + type: object + properties: + size: + type: string + enum: + - small + - medium + - large + replicas: + type: integer + minimum: 1 + maximum: 10 + required: + - size + status: + type: object + properties: + ready: + type: boolean + message: + type: string + subresources: + status: {} + additionalPrinterColumns: + - name: Size + type: string + description: The size of the widget + jsonPath: .spec.size + - name: Replicas + type: integer + description: Number of replicas + jsonPath: .spec.replicas + - name: Ready + type: boolean + description: Is the widget ready? + jsonPath: .status.ready + - name: Age + type: date + jsonPath: .metadata.creationTimestamp + diff --git a/pkg/registry/apis/apiextensions/resources/example-widget.yaml b/pkg/registry/apis/apiextensions/resources/example-widget.yaml new file mode 100644 index 00000000000..7b22240e12e --- /dev/null +++ b/pkg/registry/apis/apiextensions/resources/example-widget.yaml @@ -0,0 +1,9 @@ +apiVersion: customcrdtest.grafana.app/v1 +kind: Widget +metadata: + name: my-widget-2 + namespace: default +spec: + size: medium + replicas: 3 + diff --git a/pkg/registry/apis/apiextensions/resources/outage-crd.yaml b/pkg/registry/apis/apiextensions/resources/outage-crd.yaml new file mode 100644 index 00000000000..d50cf32b540 --- /dev/null +++ b/pkg/registry/apis/apiextensions/resources/outage-crd.yaml @@ -0,0 +1,111 @@ +apiVersion: apiextensions.k8s.io/v1 +kind: CustomResourceDefinition +metadata: + name: outages.monitoring.grafana.app +spec: + group: monitoring.grafana.app + names: + kind: Outage + listKind: OutageList + plural: outages + singular: outage + shortNames: + - out + scope: Cluster # ← CLUSTER-SCOPED (not namespaced) + versions: + - name: v1 + served: true + storage: true + schema: + openAPIV3Schema: + type: object + properties: + apiVersion: + type: string + kind: + type: string + metadata: + type: object + spec: + type: object + properties: + region: + type: string + description: Geographic region affected by the outage + enum: + - us-east-1 + - us-west-2 + - eu-west-1 + - eu-central-1 + - ap-southeast-1 + - ap-northeast-1 + severity: + type: string + description: Severity level of the outage + enum: + - critical + - major + - minor + default: major + affectedServices: + type: array + description: List of services affected + items: + type: string + startTime: + type: string + format: date-time + description: When the outage started + description: + type: string + description: Description of the outage + required: + - region + - severity + - startTime + status: + type: object + properties: + resolved: + type: boolean + description: Whether the outage has been resolved + resolvedAt: + type: string + format: date-time + description: When the outage was resolved + affectedCustomers: + type: integer + description: Number of customers affected + updates: + type: array + description: Status updates + items: + type: object + properties: + timestamp: + type: string + format: date-time + message: + type: string + subresources: + status: {} + additionalPrinterColumns: + - name: Region + type: string + description: Geographic region + jsonPath: .spec.region + - name: Severity + type: string + description: Severity level + jsonPath: .spec.severity + - name: Resolved + type: boolean + description: Resolution status + jsonPath: .status.resolved + - name: Start Time + type: date + description: When the outage started + jsonPath: .spec.startTime + - name: Age + type: date + jsonPath: .metadata.creationTimestamp diff --git a/pkg/registry/apis/apiextensions/resources/outage-instance.yaml b/pkg/registry/apis/apiextensions/resources/outage-instance.yaml new file mode 100644 index 00000000000..18d6ecf099e --- /dev/null +++ b/pkg/registry/apis/apiextensions/resources/outage-instance.yaml @@ -0,0 +1,22 @@ +apiVersion: monitoring.grafana.app/v1 +kind: Outage +metadata: + name: outage-2025-11-20-us-east + # NO namespace field - this is cluster-scoped! +spec: + region: us-east-1 + severity: critical + affectedServices: + - grafana-cloud-metrics + - grafana-cloud-logs + - grafana-cloud-traces + startTime: "2025-11-20T10:00:00Z" + description: "Database connectivity issues affecting multiple services in US East region" +status: + resolved: false + affectedCustomers: 1247 + updates: + - timestamp: "2025-11-20T10:15:00Z" + message: "Incident detected, investigating database connectivity" + - timestamp: "2025-11-20T10:30:00Z" + message: "Root cause identified, applying fix to primary database cluster" diff --git a/pkg/registry/apis/apiextensions/storage.go b/pkg/registry/apis/apiextensions/storage.go new file mode 100644 index 00000000000..54b88bd337f --- /dev/null +++ b/pkg/registry/apis/apiextensions/storage.go @@ -0,0 +1,176 @@ +package apiextensions + +import ( + "context" + "fmt" + + apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + genericregistry "k8s.io/apiserver/pkg/registry/generic/registry" + "k8s.io/apiserver/pkg/registry/rest" +) + +var _ rest.StandardStorage = (*crdStorage)(nil) +var _ rest.Scoper = (*crdStorage)(nil) + +// crdStorage wraps the generic registry store and adds CRD-specific validation and dynamic registration +type crdStorage struct { + *genericregistry.Store + dynamicReg *DynamicRegistry +} + +// NamespaceScoped returns false since CRDs are cluster-scoped resources +func (s *crdStorage) NamespaceScoped() bool { + return false +} + +// Create validates and creates a CRD, then triggers dynamic registration of the custom resource +func (s *crdStorage) Create( + ctx context.Context, + obj runtime.Object, + createValidation rest.ValidateObjectFunc, + options *metav1.CreateOptions, +) (runtime.Object, error) { + crd, ok := obj.(*apiextensionsv1.CustomResourceDefinition) + if !ok { + return nil, apierrors.NewBadRequest("object is not a CustomResourceDefinition") + } + + if err := validateCRDForPhase1(crd); err != nil { + return nil, err + } + + // Create the CRD in storage + result, err := s.Store.Create(ctx, obj, createValidation, options) + if err != nil { + return nil, err + } + + // Register the custom resource dynamically + createdCRD, ok := result.(*apiextensionsv1.CustomResourceDefinition) + if ok { + if err := s.dynamicReg.RegisterCRD(createdCRD); err != nil { + // Log error but don't fail the creation + // TODO: Add proper logging + fmt.Printf("Warning: failed to register CRD dynamically: %v\n", err) + } + } + + return result, nil +} + +// Update validates and updates a CRD, then updates the dynamic registration +func (s *crdStorage) Update( + ctx context.Context, + name string, + objInfo rest.UpdatedObjectInfo, + createValidation rest.ValidateObjectFunc, + updateValidation rest.ValidateObjectUpdateFunc, + forceAllowCreate bool, + options *metav1.UpdateOptions, +) (runtime.Object, bool, error) { + result, created, err := s.Store.Update(ctx, name, objInfo, createValidation, updateValidation, forceAllowCreate, options) + if err != nil { + return nil, false, err + } + + // Update dynamic registration + updatedCRD, ok := result.(*apiextensionsv1.CustomResourceDefinition) + if ok { + if err := s.dynamicReg.UpdateCRD(updatedCRD); err != nil { + fmt.Printf("Warning: failed to update CRD registration: %v\n", err) + } + } + + return result, created, nil +} + +// Delete deletes a CRD and unregisters the custom resource +func (s *crdStorage) Delete( + ctx context.Context, + name string, + deleteValidation rest.ValidateObjectFunc, + options *metav1.DeleteOptions, +) (runtime.Object, bool, error) { + // Get the CRD before deletion + obj, err := s.Store.Get(ctx, name, &metav1.GetOptions{}) + if err != nil { + return nil, false, err + } + + crd, ok := obj.(*apiextensionsv1.CustomResourceDefinition) + if !ok { + return nil, false, apierrors.NewInternalError(fmt.Errorf("object is not a CRD")) + } + + // Delete the CRD + result, immediate, err := s.Store.Delete(ctx, name, deleteValidation, options) + if err != nil { + return nil, false, err + } + + // Unregister the custom resource + if err := s.dynamicReg.UnregisterCRD(crd); err != nil { + fmt.Printf("Warning: failed to unregister CRD: %v\n", err) + } + + return result, immediate, nil +} + +func validateCRDForPhase1(crd *apiextensionsv1.CustomResourceDefinition) error { + // TODO(@konsalex): only support single-version CRDs for now + if len(crd.Spec.Versions) != 1 { + return apierrors.NewBadRequest( + fmt.Sprintf("we only support single-version CRDs, got %d versions", len(crd.Spec.Versions)), + ) + } + + // TODO(@konsalex): no webhook conversion for now we will need to support this + if crd.Spec.Conversion != nil && crd.Spec.Conversion.Strategy == apiextensionsv1.WebhookConverter { + return apierrors.NewBadRequest("no support webhook conversion") + } + + // Ensure the single version is marked as both served and storage + if len(crd.Spec.Versions) > 0 { + version := crd.Spec.Versions[0] + if !version.Served { + return apierrors.NewBadRequest("the single version must be marked as served") + } + if !version.Storage { + return apierrors.NewBadRequest("the single version must be marked as storage") + } + } + + return nil +} + +// crdStatusStorage handles the status subresource for CRDs +type crdStatusStorage struct { + *genericregistry.Store +} + +func (s *crdStatusStorage) New() runtime.Object { + return &apiextensionsv1.CustomResourceDefinition{} +} + +func (s *crdStatusStorage) Get( + ctx context.Context, + name string, + options *metav1.GetOptions, +) (runtime.Object, error) { + return s.Store.Get(ctx, name, options) +} + +func (s *crdStatusStorage) Update( + ctx context.Context, + name string, + objInfo rest.UpdatedObjectInfo, + createValidation rest.ValidateObjectFunc, + updateValidation rest.ValidateObjectUpdateFunc, + forceAllowCreate bool, + options *metav1.UpdateOptions, +) (runtime.Object, bool, error) { + return s.Store.Update(ctx, name, objInfo, createValidation, updateValidation, false, options) +} diff --git a/pkg/registry/apis/apiextensions/storage_cr.go b/pkg/registry/apis/apiextensions/storage_cr.go new file mode 100644 index 00000000000..dda53097fef --- /dev/null +++ b/pkg/registry/apis/apiextensions/storage_cr.go @@ -0,0 +1,687 @@ +package apiextensions + +import ( + "context" + "fmt" + + jsonpatch "github.com/evanphx/json-patch" + "github.com/google/uuid" + authlib "github.com/grafana/authlib/types" + apiextensions "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions" + apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1" + "k8s.io/apiextensions-apiserver/pkg/apiserver/validation" + apierrors "k8s.io/apimachinery/pkg/api/errors" + metainternalversion "k8s.io/apimachinery/pkg/apis/meta/internalversion" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/util/json" + "k8s.io/apimachinery/pkg/util/validation/field" + "k8s.io/apimachinery/pkg/watch" + "k8s.io/apiserver/pkg/registry/generic" + "k8s.io/apiserver/pkg/registry/rest" + "k8s.io/apiserver/pkg/storage" + + "github.com/grafana/grafana/pkg/apimachinery/utils" + grafanaregistry "github.com/grafana/grafana/pkg/apiserver/registry/generic" + "github.com/grafana/grafana/pkg/storage/unified/apistore" + "github.com/grafana/grafana/pkg/storage/unified/resource" + "github.com/grafana/grafana/pkg/storage/unified/resourcepb" +) + +var _ rest.StandardStorage = (*customResourceStorage)(nil) +var _ rest.Patcher = (*customResourceStorage)(nil) +var _ rest.Scoper = (*customResourceStorage)(nil) + +var ( + internalScheme = runtime.NewScheme() +) + +func init() { + _ = apiextensionsv1.AddToScheme(internalScheme) + _ = apiextensions.AddToScheme(internalScheme) +} + +// customResourceStorage implements REST storage for custom resource instances +// It bypasses the generic registry store and calls unified storage directly +// to properly handle unstructured objects +type customResourceStorage struct { + storage storage.Interface // Direct unified storage interface + crd *apiextensionsv1.CustomResourceDefinition + version string + gvk schema.GroupVersionKind + gvr schema.GroupVersionResource + accessClient authlib.AccessClient + keyFunc func(ctx context.Context, name string) (string, error) + keyRootFunc func(ctx context.Context) string +} + +// NewCustomResourceStorage creates storage for a custom resource based on its CRD +func NewCustomResourceStorage( + crd *apiextensionsv1.CustomResourceDefinition, + version string, + scheme *runtime.Scheme, + optsGetter generic.RESTOptionsGetter, + accessClient authlib.AccessClient, + unifiedClient resource.ResourceClient, +) (*customResourceStorage, error) { + gvk := schema.GroupVersionKind{ + Group: crd.Spec.Group, + Version: version, + Kind: crd.Spec.Names.Kind, + } + + gvr := schema.GroupVersionResource{ + Group: crd.Spec.Group, + Version: version, + Resource: crd.Spec.Names.Plural, + } + + // Register the custom resource type as unstructured + scheme.AddKnownTypeWithName(gvk, &unstructured.Unstructured{}) + listGVK := gvk + listGVK.Kind = crd.Spec.Names.ListKind + scheme.AddKnownTypeWithName(listGVK, &unstructured.UnstructuredList{}) + + // Create resource info for this custom resource + resourceInfo := utils.NewResourceInfo( + crd.Spec.Group, + version, + crd.Spec.Names.Plural, + crd.Spec.Names.Singular, + crd.Spec.Names.Kind, + func() runtime.Object { + u := &unstructured.Unstructured{} + u.SetGroupVersionKind(gvk) + return u + }, + func() runtime.Object { + ul := &unstructured.UnstructuredList{} + ul.SetGroupVersionKind(schema.GroupVersionKind{ + Group: crd.Spec.Group, + Version: version, + Kind: crd.Spec.Names.ListKind, + }) + return ul + }, + utils.TableColumns{}, + ) + + // Make it cluster-scoped if needed + if crd.Spec.Scope == apiextensionsv1.ClusterScoped { + resourceInfo = resourceInfo.WithClusterScope() + } + + // Register storage options + if restOptGetter, ok := optsGetter.(*apistore.RESTOptionsGetter); ok { + restOptGetter.RegisterOptions( + gvr.GroupResource(), + apistore.StorageOptions{}, + ) + } + + // We need this to set the codec below + gr := gvr.GroupResource() + opts, err := optsGetter.GetRESTOptions(gr, &unstructured.Unstructured{}) + if err != nil { + return nil, fmt.Errorf("failed to get REST options: %w", err) + } + + // This codec can handle any dynamic type without compile-time registration + // Without explicitly defined we get this error: + // https://github.com/kubernetes/apimachinery/blob/5a348c53eef0072c40ddf00a45ace423c2790f2a/pkg/runtime/error.go#L52 + opts.StorageConfig.Codec = unstructured.UnstructuredJSONScheme + + // Get the ConfigForResource + config := opts.StorageConfig.ForResource(gr) + + // Create key functions for storage + keyFunc := func(obj runtime.Object) (string, error) { + accessor, err := utils.MetaAccessor(obj) + if err != nil { + return "", err + } + name := accessor.GetName() + ns := accessor.GetNamespace() + + key := &grafanaregistry.Key{ + Group: gr.Group, + Resource: gr.Resource, + Namespace: ns, + Name: name, + } + return key.String(), nil + } + + keyParser := func(key string) (*resourcepb.ResourceKey, error) { + k, err := grafanaregistry.ParseKey(key) + if err != nil { + return nil, err + } + return &resourcepb.ResourceKey{ + Namespace: k.Namespace, + Group: k.Group, + Resource: k.Resource, + Name: k.Name, + }, nil + } + + // Create the actual storage using apistore.NewStorage + underlyingStorage, _, err := apistore.NewStorage( + config, + unifiedClient, + keyFunc, + keyParser, + func() runtime.Object { return &unstructured.Unstructured{} }, + func() runtime.Object { return &unstructured.UnstructuredList{} }, + grafanaregistry.GetAttrs, + nil, // trigger + nil, // indexers + nil, // configProvider + apistore.StorageOptions{}, + ) + if err != nil { + return nil, fmt.Errorf("failed to create storage: %w", err) + } + + return &customResourceStorage{ + storage: underlyingStorage, + crd: crd, + version: version, + gvk: gvk, + gvr: gvr, + accessClient: accessClient, + keyFunc: grafanaregistry.NamespaceKeyFunc(gr), + keyRootFunc: grafanaregistry.KeyRootFunc(gr), + }, nil +} + +// NamespaceScoped returns whether this custom resource is namespaced +func (s *customResourceStorage) NamespaceScoped() bool { + return s.crd.Spec.Scope == apiextensionsv1.NamespaceScoped +} + +// New returns a new instance of the custom resource +func (s *customResourceStorage) New() runtime.Object { + return &unstructured.Unstructured{} +} + +// NewList returns a new list instance +func (s *customResourceStorage) NewList() runtime.Object { + return &unstructured.UnstructuredList{} +} + +// getValidator returns a validator for the custom resource +func (s *customResourceStorage) getValidator() (validation.SchemaValidator, error) { + // Find the version we are serving + var version *apiextensionsv1.CustomResourceDefinitionVersion + for i := range s.crd.Spec.Versions { + v := &s.crd.Spec.Versions[i] + if v.Name == s.version { + version = v + break + } + } + + // If no schema is defined, we can't validate + if version == nil || version.Schema == nil || version.Schema.OpenAPIV3Schema == nil { + return nil, nil + } + + // Convert v1 schema to internal schema + internalSchema := &apiextensions.JSONSchemaProps{} + if err := internalScheme.Convert(version.Schema.OpenAPIV3Schema, internalSchema, nil); err != nil { + return nil, fmt.Errorf("failed to convert schema to internal version: %w", err) + } + + // Create validator + val, _, err := validation.NewSchemaValidator(internalSchema) + return val, err +} + +// Create validates and creates a custom resource instance +func (s *customResourceStorage) Create( + ctx context.Context, + obj runtime.Object, + createValidation rest.ValidateObjectFunc, + options *metav1.CreateOptions, +) (runtime.Object, error) { + // Ensure the object is unstructured as it is a "custom" CRD + u, ok := obj.(*unstructured.Unstructured) + if !ok { + return nil, apierrors.NewBadRequest("object must be unstructured") + } + + // Validate GVK matches the CRD + // Could it arrive here without a match from storage match failure? + objGVK := u.GroupVersionKind() + if objGVK.Group != s.gvk.Group || objGVK.Version != s.gvk.Version || objGVK.Kind != s.gvk.Kind { + return nil, apierrors.NewBadRequest( + fmt.Sprintf("object GVK %s does not match CRD GVK %s", objGVK.String(), s.gvk.String()), + ) + } + + // Set defaults if not present + if u.GetNamespace() == "" && s.NamespaceScoped() { + u.SetNamespace("default") + } + if u.GetName() == "" { + u.SetName(uuid.New().String()) + } + + // Ensure GVK is set + u.SetGroupVersionKind(s.gvk) + + // Run validation if provided + if createValidation != nil { + if err := createValidation(ctx, obj); err != nil { + return nil, err + } + } + + // Validate against OpenAPI schema + validator, err := s.getValidator() + if err != nil { + return nil, apierrors.NewInternalError(fmt.Errorf("failed to get validator: %v", err)) + } + if validator != nil { + if errs := validation.ValidateCustomResource(field.NewPath(""), u.UnstructuredContent(), validator); len(errs) > 0 { + return nil, apierrors.NewInvalid(s.gvk.GroupKind(), u.GetName(), errs) + } + } + + // Generate storage key + key, err := s.keyFunc(ctx, u.GetName()) + if err != nil { + return nil, fmt.Errorf("failed to generate key: %w", err) + } + + if options != nil && len(options.DryRun) > 0 { + return u, nil + } + + // Create the object in storage + // TODO(@konsalex): Figure out if this ttl should be (!=0) + out := &unstructured.Unstructured{} + if err := s.storage.Create(ctx, key, obj, out, 0); err != nil { + return nil, err + } + + return out, nil +} + +// Get retrieves a custom resource instance by name +func (s *customResourceStorage) Get( + ctx context.Context, + name string, + options *metav1.GetOptions, +) (runtime.Object, error) { + key, err := s.keyFunc(ctx, name) + if err != nil { + return nil, err + } + + out := &unstructured.Unstructured{} + getOpts := storage.GetOptions{ + IgnoreNotFound: false, + ResourceVersion: "", + } + if options != nil { + getOpts.ResourceVersion = options.ResourceVersion + } + if err := s.storage.Get(ctx, key, getOpts, out); err != nil { + return nil, err + } + + return out, nil +} + +// List retrieves a list of custom resource instances +func (s *customResourceStorage) List( + ctx context.Context, + options *metainternalversion.ListOptions, +) (runtime.Object, error) { + // Create an empty list to populate + listObj := &unstructured.UnstructuredList{} + listObj.SetGroupVersionKind(schema.GroupVersionKind{ + Group: s.gvk.Group, + Version: s.gvk.Version, + Kind: s.crd.Spec.Names.ListKind, + }) + + // Handle nil options + if options == nil { + options = &metainternalversion.ListOptions{} + } + + // Convert metainternalversion.ListOptions to storage.ListOptions + listOpts := storage.ListOptions{ + ResourceVersion: options.ResourceVersion, + ResourceVersionMatch: options.ResourceVersionMatch, + Predicate: storage.Everything, + } + + if options.LabelSelector != nil { + listOpts.Predicate.Label = options.LabelSelector + } + if options.FieldSelector != nil { + listOpts.Predicate.Field = options.FieldSelector + } + + // Get the key prefix for listing + keyPrefix := s.keyRootFunc(ctx) + + // Call the underlying storage List/GetList + if err := s.storage.GetList(ctx, keyPrefix, listOpts, listObj); err != nil { + return nil, err + } + + return listObj, nil +} + +// Delete removes a custom resource instance +func (s *customResourceStorage) Delete( + ctx context.Context, + name string, + deleteValidation rest.ValidateObjectFunc, + options *metav1.DeleteOptions, +) (runtime.Object, bool, error) { + key, err := s.keyFunc(ctx, name) + if err != nil { + return nil, false, err + } + + out := &unstructured.Unstructured{} + preconditions := &storage.Preconditions{} + if options != nil && options.Preconditions != nil { + if options.Preconditions.UID != nil { + preconditions.UID = options.Preconditions.UID + } + if options.Preconditions.ResourceVersion != nil { + preconditions.ResourceVersion = options.Preconditions.ResourceVersion + } + } + + validateFunc := func(ctx context.Context, obj runtime.Object) error { + if deleteValidation != nil { + return deleteValidation(ctx, obj) + } + return nil + } + + if options != nil && len(options.DryRun) > 0 { + // We need to fetch it to make sure it exists and to return it + if err := s.storage.Get(ctx, key, storage.GetOptions{}, out); err != nil { + return nil, false, err + } + if err := validateFunc(ctx, out); err != nil { + return nil, false, err + } + return out, true, nil + } + + if err := s.storage.Delete(ctx, key, out, preconditions, validateFunc, out, storage.DeleteOptions{}); err != nil { + return nil, false, err + } + + return out, true, nil +} + +// Update updates a custom resource instance +func (s *customResourceStorage) Update( + ctx context.Context, + name string, + objInfo rest.UpdatedObjectInfo, + createValidation rest.ValidateObjectFunc, + updateValidation rest.ValidateObjectUpdateFunc, + forceAllowCreate bool, + options *metav1.UpdateOptions, +) (runtime.Object, bool, error) { + key, err := s.keyFunc(ctx, name) + if err != nil { + return nil, false, err + } + + // Get the existing object + existingObj := &unstructured.Unstructured{} + existingObj.SetGroupVersionKind(s.gvk) + + err = s.storage.Get(ctx, key, storage.GetOptions{}, existingObj) + if err != nil { + if storage.IsNotFound(err) { + if !forceAllowCreate { + return nil, false, apierrors.NewNotFound(s.gvr.GroupResource(), name) + } + // If forceAllowCreate is true, treat as a create + // and not as an update + newObj, err := objInfo.UpdatedObject(ctx, nil) + if err != nil { + return nil, false, err + } + + if createValidation != nil { + if err := createValidation(ctx, newObj); err != nil { + return nil, false, err + } + } + + // Call Create internally + created, err := s.Create(ctx, newObj, nil, nil) + if err != nil { + return nil, false, err + } + return created, true, nil + } + return nil, false, err + } + + // Get the updated object + updatedObj, err := objInfo.UpdatedObject(ctx, existingObj) + if err != nil { + return nil, false, err + } + + // Validate the update + if updateValidation != nil { + if err := updateValidation(ctx, updatedObj, existingObj); err != nil { + return nil, false, err + } + } + + // Ensure it's an unstructured object + updatedUnstructured, ok := updatedObj.(*unstructured.Unstructured) + if !ok { + return nil, false, fmt.Errorf("updated object is not unstructured") + } + + // Set the GVK + updatedUnstructured.SetGroupVersionKind(s.gvk) + + // Ensure namespace and name are set correctly + if s.NamespaceScoped() { + ns := updatedUnstructured.GetNamespace() + if ns == "" { + updatedUnstructured.SetNamespace("default") + } + } + updatedUnstructured.SetName(name) + + // Validate against OpenAPI schema + validator, err := s.getValidator() + if err != nil { + return nil, false, apierrors.NewInternalError(fmt.Errorf("failed to get validator: %v", err)) + } + if validator != nil { + // We need the old object as interface{} + if errs := validation.ValidateCustomResourceUpdate(field.NewPath(""), updatedUnstructured.UnstructuredContent(), existingObj.UnstructuredContent(), validator); len(errs) > 0 { + return nil, false, apierrors.NewInvalid(s.gvk.GroupKind(), name, errs) + } + } + + // Use GuaranteedUpdate for optimistic concurrency control + out := &unstructured.Unstructured{} + out.SetGroupVersionKind(s.gvk) + + updateFunc := func(input runtime.Object, respMeta storage.ResponseMeta) (runtime.Object, *uint64, error) { + // Return the updated object + return updatedUnstructured, nil, nil + } + + if options != nil && len(options.DryRun) > 0 { + return updatedUnstructured, false, nil + } + + preconditions := &storage.Preconditions{} + + // We explicitly ignore the not-found, as we checked before proceeding + err = s.storage.GuaranteedUpdate(ctx, key, out, true, preconditions, updateFunc, out) + if err != nil { + return nil, false, err + } + + return out, false, nil +} + +// Patch patches a custom resource instance +func (s *customResourceStorage) Patch( + ctx context.Context, + name string, + patchType types.PatchType, + patchBytes []byte, + options *metav1.PatchOptions, + subresources ...string, +) (runtime.Object, error) { + key, err := s.keyFunc(ctx, name) + if err != nil { + return nil, err + } + + existingObj := &unstructured.Unstructured{} + existingObj.SetGroupVersionKind(s.gvk) + + err = s.storage.Get(ctx, key, storage.GetOptions{}, existingObj) + if err != nil { + if storage.IsNotFound(err) { + return nil, apierrors.NewNotFound(s.gvr.GroupResource(), name) + } + return nil, err + } + + existingJSON, err := json.Marshal(existingObj.Object) + if err != nil { + return nil, fmt.Errorf("failed to marshal existing object: %w", err) + } + + var patchedJSON []byte + + switch patchType { + case types.JSONPatchType: + patch, err := jsonpatch.DecodePatch(patchBytes) + if err != nil { + return nil, fmt.Errorf("failed to decode JSON patch: %w", err) + } + patchedJSON, err = patch.Apply(existingJSON) + if err != nil { + return nil, fmt.Errorf("failed to apply JSON patch: %w", err) + } + + case types.MergePatchType: + patchedJSON, err = jsonpatch.MergePatch(existingJSON, patchBytes) + if err != nil { + return nil, fmt.Errorf("failed to apply merge patch: %w", err) + } + + case types.StrategicMergePatchType: + // Kubernetes Strategic Merge Patch + // For unstructured objects, strategic merge patch behaves like merge patch + // because we don't have a Go struct to define merge strategies + // TODO(@konsalex): Clarify is we need to even support this by falling-back to merge patch, or just return an error to inform clients + // patchedJSON, err = strategicpatch.StrategicMergePatch(existingJSON, patchBytes, &unstructured.Unstructured{}) + + // Fallback to merge patch if strategic merge fails + patchedJSON, err = jsonpatch.MergePatch(existingJSON, patchBytes) + if err != nil { + return nil, fmt.Errorf("failed to apply patch: %w", err) + } + + default: + return nil, fmt.Errorf("unsupported patch type: %s", patchType) + } + + // Unmarshal the patched JSON into an unstructured object + patchedObj := &unstructured.Unstructured{} + if err := json.Unmarshal(patchedJSON, &patchedObj.Object); err != nil { + return nil, fmt.Errorf("failed to unmarshal patched object: %w", err) + } + + // Set the GVK + patchedObj.SetGroupVersionKind(s.gvk) + + // Ensure namespace and name are set correctly + if s.NamespaceScoped() { + ns := patchedObj.GetNamespace() + if ns == "" { + patchedObj.SetNamespace(existingObj.GetNamespace()) + } + } + patchedObj.SetName(name) + + // Validate against OpenAPI schema + validator, err := s.getValidator() + if err != nil { + return nil, apierrors.NewInternalError(fmt.Errorf("failed to get validator: %v", err)) + } + if validator != nil { + // We need the old object as interface{} + if errs := validation.ValidateCustomResourceUpdate(field.NewPath(""), patchedObj.UnstructuredContent(), existingObj.UnstructuredContent(), validator); len(errs) > 0 { + return nil, apierrors.NewInvalid(s.gvk.GroupKind(), name, errs) + } + } + + // Use GuaranteedUpdate to save the patched object + out := &unstructured.Unstructured{} + out.SetGroupVersionKind(s.gvk) + + updateFunc := func(input runtime.Object, respMeta storage.ResponseMeta) (runtime.Object, *uint64, error) { + return patchedObj, nil, nil + } + + preconditions := &storage.Preconditions{} + if options != nil && options.DryRun != nil && len(options.DryRun) > 0 { + fmt.Printf(" - Dry run mode\n") + } + + err = s.storage.GuaranteedUpdate(ctx, key, out, true, preconditions, updateFunc, out) + if err != nil { + return nil, err + } + + return out, nil +} + +// DeleteCollection deletes a collection of custom resources +func (s *customResourceStorage) DeleteCollection( + ctx context.Context, + deleteValidation rest.ValidateObjectFunc, + options *metav1.DeleteOptions, + listOptions *metainternalversion.ListOptions, +) (runtime.Object, error) { + return nil, apierrors.NewMethodNotSupported(s.gvr.GroupResource(), "deletecollection") +} + +// Watch returns a watch interface for custom resources +func (s *customResourceStorage) Watch(ctx context.Context, options *metainternalversion.ListOptions) (watch.Interface, error) { + return nil, apierrors.NewMethodNotSupported(s.gvr.GroupResource(), "watch") +} + +// ConvertToTable converts to a table for kubectl +func (s *customResourceStorage) ConvertToTable(ctx context.Context, object runtime.Object, tableOptions runtime.Object) (*metav1.Table, error) { + return nil, apierrors.NewMethodNotSupported(s.gvr.GroupResource(), "table") +} + +// Destroy cleans up resources +func (s *customResourceStorage) Destroy() { + // Nothing to clean up for direct storage +} diff --git a/pkg/registry/apis/apis.go b/pkg/registry/apis/apis.go index 10f1d7e52be..15aa03548de 100644 --- a/pkg/registry/apis/apis.go +++ b/pkg/registry/apis/apis.go @@ -1,6 +1,7 @@ package apiregistry import ( + "github.com/grafana/grafana/pkg/registry/apis/apiextensions" "github.com/grafana/grafana/pkg/registry/apis/collections" dashboardinternal "github.com/grafana/grafana/pkg/registry/apis/dashboard" "github.com/grafana/grafana/pkg/registry/apis/dashboardsnapshot" @@ -20,6 +21,7 @@ type Service struct{} // ProvideRegistryServiceSink is an entry point for each service that will force initialization // and give each builder the chance to register itself with the main server func ProvideRegistryServiceSink( + _ *apiextensions.APIExtensionsBuilder, _ *dashboardinternal.DashboardsAPIBuilder, _ *dashboardsnapshot.SnapshotsAPIBuilder, _ *datasource.DataSourceAPIBuilder, diff --git a/pkg/registry/apis/wireset.go b/pkg/registry/apis/wireset.go index d19b405c6c9..1d542d6cbff 100644 --- a/pkg/registry/apis/wireset.go +++ b/pkg/registry/apis/wireset.go @@ -3,6 +3,7 @@ package apiregistry import ( "github.com/google/wire" + "github.com/grafana/grafana/pkg/registry/apis/apiextensions" "github.com/grafana/grafana/pkg/registry/apis/collections" dashboardinternal "github.com/grafana/grafana/pkg/registry/apis/dashboard" "github.com/grafana/grafana/pkg/registry/apis/dashboardsnapshot" @@ -59,6 +60,7 @@ var WireSet = wire.NewSet( provisioningExtras, // Each must be added here *and* in the ServiceSink above + apiextensions.RegisterAPIService, dashboardinternal.RegisterAPIService, dashboardsnapshot.RegisterAPIService, datasource.RegisterAPIService, diff --git a/pkg/server/wire_gen.go b/pkg/server/wire_gen.go index 279a67f4134..3de4c344e44 100644 --- a/pkg/server/wire_gen.go +++ b/pkg/server/wire_gen.go @@ -48,6 +48,7 @@ import ( "github.com/grafana/grafana/pkg/plugins/pluginscdn" "github.com/grafana/grafana/pkg/plugins/repo" "github.com/grafana/grafana/pkg/registry/apis" + "github.com/grafana/grafana/pkg/registry/apis/apiextensions" "github.com/grafana/grafana/pkg/registry/apis/collections" "github.com/grafana/grafana/pkg/registry/apis/dashboard" "github.com/grafana/grafana/pkg/registry/apis/dashboard/legacy" @@ -861,6 +862,10 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api identitySynchronizer := authnimpl.ProvideIdentitySynchronizer(authnimplService) ldapImpl := service12.ProvideService(cfg, featureToggles, ssosettingsimplService) apiService := api4.ProvideService(cfg, routeRegisterImpl, accessControl, userService, authinfoimplService, ossGroups, identitySynchronizer, orgService, ldapImpl, userAuthTokenService, bundleregistryService) + apiExtensionsBuilder, err := apiextensions.RegisterAPIService(cfg, featureToggles, apiserverService, accessClient, registerer, resourceClient) + if err != nil { + return nil, err + } dashboardsAPIBuilder := dashboard.RegisterAPIService(cfg, featureToggles, apiserverService, dashboardService, dashboardProvisioningService, service15, dashboardServiceImpl, dashboardPermissionsService, accessControl, accessClient, provisioningServiceImpl, dashboardsStore, registerer, sqlStore, tracingService, resourceClient, dualwriteService, sortService, quotaService, libraryPanelService, eventualRestConfigProvider, userService, libraryElementService, publicDashboardServiceImpl) snapshotsAPIBuilder := dashboardsnapshot.RegisterAPIService(serviceImpl, apiserverService, cfg, featureToggles, sqlStore, registerer) dataSourceAPIBuilder, err := datasource.RegisterAPIService(featureToggles, apiserverService, middlewareHandler, scopedPluginDatasourceProvider, plugincontextProvider, accessControl, registerer, sourcesService) @@ -914,7 +919,7 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api if err != nil { return nil, err } - apiregistryService := apiregistry.ProvideRegistryServiceSink(dashboardsAPIBuilder, snapshotsAPIBuilder, dataSourceAPIBuilder, folderAPIBuilder, identityAccessManagementAPIBuilder, queryAPIBuilder, userStorageAPIBuilder, apiBuilder, collectionsAPIBuilder, provisioningAPIBuilder, ofrepAPIBuilder, dependencyRegisterer, provisioningDependencyRegisterer) + apiregistryService := apiregistry.ProvideRegistryServiceSink(apiExtensionsBuilder, dashboardsAPIBuilder, snapshotsAPIBuilder, dataSourceAPIBuilder, folderAPIBuilder, identityAccessManagementAPIBuilder, queryAPIBuilder, userStorageAPIBuilder, apiBuilder, collectionsAPIBuilder, provisioningAPIBuilder, ofrepAPIBuilder, dependencyRegisterer, provisioningDependencyRegisterer) teamPermissionsService, err := ossaccesscontrol.ProvideTeamPermissions(cfg, featureToggles, routeRegisterImpl, sqlStore, accessControl, ossLicensingService, acimplService, teamService, userService, actionSetService) if err != nil { return nil, err @@ -1511,6 +1516,10 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac identitySynchronizer := authnimpl.ProvideIdentitySynchronizer(authnimplService) ldapImpl := service12.ProvideService(cfg, featureToggles, ssosettingsimplService) apiService := api4.ProvideService(cfg, routeRegisterImpl, accessControl, userService, authinfoimplService, ossGroups, identitySynchronizer, orgService, ldapImpl, userAuthTokenService, bundleregistryService) + apiExtensionsBuilder, err := apiextensions.RegisterAPIService(cfg, featureToggles, apiserverService, accessClient, registerer, resourceClient) + if err != nil { + return nil, err + } dashboardsAPIBuilder := dashboard.RegisterAPIService(cfg, featureToggles, apiserverService, dashboardService, dashboardProvisioningService, service15, dashboardServiceImpl, dashboardPermissionsService, accessControl, accessClient, provisioningServiceImpl, dashboardsStore, registerer, sqlStore, tracingService, resourceClient, dualwriteService, sortService, quotaService, libraryPanelService, eventualRestConfigProvider, userService, libraryElementService, publicDashboardServiceImpl) snapshotsAPIBuilder := dashboardsnapshot.RegisterAPIService(serviceImpl, apiserverService, cfg, featureToggles, sqlStore, registerer) dataSourceAPIBuilder, err := datasource.RegisterAPIService(featureToggles, apiserverService, middlewareHandler, scopedPluginDatasourceProvider, plugincontextProvider, accessControl, registerer, sourcesService) @@ -1564,7 +1573,7 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac if err != nil { return nil, err } - apiregistryService := apiregistry.ProvideRegistryServiceSink(dashboardsAPIBuilder, snapshotsAPIBuilder, dataSourceAPIBuilder, folderAPIBuilder, identityAccessManagementAPIBuilder, queryAPIBuilder, userStorageAPIBuilder, apiBuilder, collectionsAPIBuilder, provisioningAPIBuilder, ofrepAPIBuilder, dependencyRegisterer, provisioningDependencyRegisterer) + apiregistryService := apiregistry.ProvideRegistryServiceSink(apiExtensionsBuilder, dashboardsAPIBuilder, snapshotsAPIBuilder, dataSourceAPIBuilder, folderAPIBuilder, identityAccessManagementAPIBuilder, queryAPIBuilder, userStorageAPIBuilder, apiBuilder, collectionsAPIBuilder, provisioningAPIBuilder, ofrepAPIBuilder, dependencyRegisterer, provisioningDependencyRegisterer) teamPermissionsService, err := ossaccesscontrol.ProvideTeamPermissions(cfg, featureToggles, routeRegisterImpl, sqlStore, accessControl, ossLicensingService, acimplService, teamService, userService, actionSetService) if err != nil { return nil, err diff --git a/pkg/services/apiserver/service.go b/pkg/services/apiserver/service.go index 605016c112d..d6235c02643 100644 --- a/pkg/services/apiserver/service.go +++ b/pkg/services/apiserver/service.go @@ -379,16 +379,29 @@ func (s *service) start(ctx context.Context) error { notFoundHandler := notfoundhandler.New(s.codecs, genericapifilters.NoMuxAndDiscoveryIncompleteKey) + // Wrap the not-found handler with dynamic custom resource handler + finalHandler := s.createDynamicHandlerWrapper(builders, notFoundHandler) + if err := appinstaller.RegisterPostStartHooks(s.appInstallers, serverConfig); err != nil { return fmt.Errorf("failed to register post start hooks for app installers: %w", err) } // Create the server - server, err := serverConfig.Complete().New("grafana-apiserver", genericapiserver.NewEmptyDelegateWithCustomHandler(notFoundHandler)) + server, err := serverConfig.Complete().New("grafana-apiserver", genericapiserver.NewEmptyDelegateWithCustomHandler(finalHandler)) if err != nil { return err } + // Inject the server instance into any builders that need it (e.g., APIExtensionsBuilder for dynamic CRD registration) + type apiServerSetter interface { + SetAPIServer(server *genericapiserver.GenericAPIServer) + } + for _, b := range builders { + if setter, ok := b.(apiServerSetter); ok { + setter.SetAPIServer(server) + } + } + // Install the API group+version for existing builders err = builder.InstallAPIs(s.scheme, s.codecs, @@ -497,6 +510,77 @@ func (s *service) start(ctx context.Context) error { return nil } +// createDynamicHandlerWrapper wraps the not-found handler with dynamic custom resource handlers +func (s *service) createDynamicHandlerWrapper(builders []builder.APIGroupBuilder, notFoundHandler http.Handler) http.Handler { + // Look for the APIExtensionsBuilder - we'll fetch the handler lazily on each request + // because the handler is created during UpdateAPIGroupInfo which happens AFTER this wrapper is installed + fmt.Printf("createDynamicHandlerWrapper: Setting up lazy handler wrapper for %d builders...\n", len(builders)) + + type dynamicHandlerProvider interface { + GetDynamicHandler() http.Handler + } + + var dhProvider dynamicHandlerProvider + for i, b := range builders { + fmt.Printf(" Builder %d: %T\n", i, b) + if provider, ok := b.(dynamicHandlerProvider); ok { + fmt.Printf(" ✓ Builder %d implements dynamicHandlerProvider - will fetch handler lazily\n", i) + dhProvider = provider + break + } + } + + if dhProvider == nil { + // No dynamic handler provider found, just use the not-found handler + fmt.Println(" No dynamic handler provider found, using standard not-found handler") + return notFoundHandler + } + + // Return a wrapper that lazily fetches and tries the dynamic handler on each request + return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + // Check if this is an /apis/ request that might be for a custom resource + if strings.HasPrefix(r.URL.Path, "/apis/") { + // Fetch the dynamic handler (it will be nil until UpdateAPIGroupInfo creates it) + dynamicHandler := dhProvider.GetDynamicHandler() + if dynamicHandler != nil { + // Create a response recorder to capture what the dynamic handler does + recorder := &responseRecorder{ + ResponseWriter: w, + statusCode: 0, + } + + dynamicHandler.ServeHTTP(recorder, r) + + // If the dynamic handler handled it (didn't return 404), we're done + if recorder.statusCode != 0 && recorder.statusCode != http.StatusNotFound { + return + } + } + } + + // Otherwise, fall back to the not-found handler + notFoundHandler.ServeHTTP(w, r) + }) +} + +// responseRecorder captures the status code from a handler +type responseRecorder struct { + http.ResponseWriter + statusCode int +} + +func (r *responseRecorder) WriteHeader(statusCode int) { + r.statusCode = statusCode + r.ResponseWriter.WriteHeader(statusCode) +} + +func (r *responseRecorder) Write(b []byte) (int, error) { + if r.statusCode == 0 { + r.statusCode = http.StatusOK + } + return r.ResponseWriter.Write(b) +} + func (s *service) startCoreServer( ctx context.Context, transport *grafanaapiserveroptions.RoundTripperFunc, diff --git a/pkg/services/featuremgmt/registry.go b/pkg/services/featuremgmt/registry.go index 65f94c24eac..5b4d381f4e4 100644 --- a/pkg/services/featuremgmt/registry.go +++ b/pkg/services/featuremgmt/registry.go @@ -755,6 +755,13 @@ var ( Owner: grafanaAppPlatformSquad, RequiresRestart: true, }, + { + Name: "apiExtensions", + Description: "Enable Kubernetes CustomResourceDefinition (CRD) support with dynamic API registration", + Stage: FeatureStageExperimental, + Owner: grafanaAppPlatformSquad, + RequiresRestart: true, + }, { Name: "groupByVariable", Description: "Enable groupBy variable support in scenes dashboards", diff --git a/pkg/services/featuremgmt/toggles_gen.go b/pkg/services/featuremgmt/toggles_gen.go index b88204db708..7c5a74080e6 100644 --- a/pkg/services/featuremgmt/toggles_gen.go +++ b/pkg/services/featuremgmt/toggles_gen.go @@ -311,6 +311,10 @@ const ( // Enable CAP token based authentication in grafana's embedded kube-aggregator FlagKubernetesAggregatorCapTokenAuth = "kubernetesAggregatorCapTokenAuth" + // FlagApiExtensions + // Enable Kubernetes CustomResourceDefinition (CRD) support with dynamic API registration + FlagApiExtensions = "apiExtensions" + // FlagGroupByVariable // Enable groupBy variable support in scenes dashboards FlagGroupByVariable = "groupByVariable" diff --git a/pkg/storage/unified/resource/access.go b/pkg/storage/unified/resource/access.go index 58a9d8b8b76..cafa972e7a3 100644 --- a/pkg/storage/unified/resource/access.go +++ b/pkg/storage/unified/resource/access.go @@ -135,7 +135,8 @@ func (c authzLimitedClient) Check(ctx context.Context, id claims.AuthInfo, req c return claims.CheckResponse{Allowed: true}, nil } - if !claims.NamespaceMatches(id.GetNamespace(), req.Namespace) { + // For cluster-scoped resources (empty namespace), skip namespace matching check + if req.Namespace != "" && !claims.NamespaceMatches(id.GetNamespace(), req.Namespace) { span.SetAttributes(attribute.Bool("allowed", false)) span.SetStatus(codes.Error, "Namespace mismatch") span.RecordError(claims.ErrNamespaceMismatch) @@ -184,7 +185,8 @@ func (c authzLimitedClient) Compile(ctx context.Context, id claims.AuthInfo, req return true }, claims.NoopZookie{}, nil } - if !claims.NamespaceMatches(id.GetNamespace(), req.Namespace) { + // For cluster-scoped resources (empty namespace), skip namespace matching check + if req.Namespace != "" && !claims.NamespaceMatches(id.GetNamespace(), req.Namespace) { span.SetAttributes(attribute.Bool("allowed", false)) span.SetStatus(codes.Error, "Namespace mismatch") span.RecordError(claims.ErrNamespaceMismatch)