Revert "Revert: DataSource: Support config CRUD from apiservers (#106996) (#110342)"

This reverts commit 72eeefabd7.
This commit is contained in:
beejeebus
2025-10-07 14:31:07 -04:00
committed by beejeebus
parent 84a2f41016
commit c3f34efb41
58 changed files with 2039 additions and 471 deletions
+161
View File
@@ -0,0 +1,161 @@
package query
import (
"context"
"fmt"
"k8s.io/apimachinery/pkg/apis/meta/internalversion"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apiserver/pkg/endpoints/request"
"k8s.io/apiserver/pkg/registry/rest"
authlib "github.com/grafana/authlib/types"
"github.com/grafana/grafana/pkg/apimachinery/utils"
queryV0 "github.com/grafana/grafana/pkg/apis/query/v0alpha1"
gapiutil "github.com/grafana/grafana/pkg/services/apiserver/utils"
"github.com/grafana/grafana/pkg/services/datasources"
)
var (
_ rest.Scoper = (*connectionAccess)(nil)
_ rest.SingularNameProvider = (*connectionAccess)(nil)
_ rest.Getter = (*connectionAccess)(nil)
_ rest.Lister = (*connectionAccess)(nil)
_ rest.Storage = (*connectionAccess)(nil)
)
// Get all datasource connections -- this will be backed by search or duplicated resource in unified storage
type DataSourceConnectionProvider interface {
// Get gets a specific datasource (that the user in context can see)
// The name is {group}:{name}, see /pkg/apis/query/v0alpha1/connection.go#L34
GetConnection(ctx context.Context, namespace string, name string) (*queryV0.DataSourceConnection, error)
// List lists all data sources the user in context can see
ListConnections(ctx context.Context, namespace string) (*queryV0.DataSourceConnectionList, error)
}
type connectionAccess struct {
tableConverter rest.TableConvertor
connections DataSourceConnectionProvider
}
func (s *connectionAccess) New() runtime.Object {
return queryV0.ConnectionResourceInfo.NewFunc()
}
func (s *connectionAccess) Destroy() {}
func (s *connectionAccess) NamespaceScoped() bool {
return true
}
func (s *connectionAccess) GetSingularName() string {
return queryV0.ConnectionResourceInfo.GetSingularName()
}
func (s *connectionAccess) ShortNames() []string {
return queryV0.ConnectionResourceInfo.GetShortNames()
}
func (s *connectionAccess) NewList() runtime.Object {
return queryV0.ConnectionResourceInfo.NewListFunc()
}
func (s *connectionAccess) ConvertToTable(ctx context.Context, object runtime.Object, tableOptions runtime.Object) (*metav1.Table, error) {
if s.tableConverter == nil {
s.tableConverter = queryV0.ConnectionResourceInfo.TableConverter()
}
return s.tableConverter.ConvertToTable(ctx, object, tableOptions)
}
func (s *connectionAccess) Get(ctx context.Context, name string, options *metav1.GetOptions) (runtime.Object, error) {
return s.connections.GetConnection(ctx, request.NamespaceValue(ctx), name)
}
func (s *connectionAccess) List(ctx context.Context, options *internalversion.ListOptions) (runtime.Object, error) {
return s.connections.ListConnections(ctx, request.NamespaceValue(ctx))
}
type connectionsProvider struct {
dsService datasources.DataSourceService
registry queryV0.DataSourceApiServerRegistry
}
var (
_ DataSourceConnectionProvider = (*connectionsProvider)(nil)
)
func (q *connectionsProvider) GetConnection(ctx context.Context, namespace string, name string) (*queryV0.DataSourceConnection, error) {
info, err := authlib.ParseNamespace(namespace)
if err != nil {
return nil, err
}
ds, err := q.dsService.GetDataSource(ctx, &datasources.GetDataSourceQuery{
UID: name,
OrgID: info.OrgID,
})
if err != nil {
return nil, err
}
// TODO... access control?
return q.asConnection(ds, namespace)
}
func (q *connectionsProvider) ListConnections(ctx context.Context, namespace string) (*queryV0.DataSourceConnectionList, error) {
ns, err := authlib.ParseNamespace(namespace)
if err != nil {
return nil, err
}
dss, err := q.dsService.GetDataSources(ctx, &datasources.GetDataSourcesQuery{
OrgID: ns.OrgID,
DataSourceLimit: 10000,
})
if err != nil {
return nil, err
}
result := &queryV0.DataSourceConnectionList{
Items: []queryV0.DataSourceConnection{},
}
for _, ds := range dss {
v, err := q.asConnection(ds, namespace)
if err != nil {
return nil, err
}
result.Items = append(result.Items, *v)
}
return result, nil
}
func (q *connectionsProvider) asConnection(ds *datasources.DataSource, ns string) (v *queryV0.DataSourceConnection, err error) {
gv, err := q.registry.GetDatasourceGroupVersion(ds.Type)
if err != nil {
return nil, fmt.Errorf("datasource type %q does not map to an apiserver %w", ds.Type, err)
}
v = &queryV0.DataSourceConnection{
ObjectMeta: metav1.ObjectMeta{
Name: queryV0.DataSourceConnectionName(gv.Group, ds.UID),
Namespace: ns,
CreationTimestamp: metav1.NewTime(ds.Created),
ResourceVersion: fmt.Sprintf("%d", ds.Updated.UnixMilli()),
Generation: int64(ds.Version),
},
Title: ds.Name,
Datasource: queryV0.DataSourceConnectionRef{
Group: gv.Group,
Version: gv.Version,
Name: ds.UID,
},
}
v.UID = gapiutil.CalculateClusterWideUID(v) // UID is unique across all groups
if !ds.Updated.IsZero() {
meta, err := utils.MetaAccessor(v)
if err != nil {
meta.SetUpdatedTimestamp(&ds.Updated)
}
}
return v, err
}
@@ -66,21 +66,8 @@ func AddQueriesToOpenAPI(options OASQueryOptions) error {
// Rewrite the query path
query := oas.Paths.Paths[root+options.QueryPath]
if query != nil && query.Post != nil {
query.Post.Tags = []string{"Query"}
query.Parameters = []*spec3.Parameter{
{
ParameterProps: spec3.ParameterProps{
Name: "namespace",
In: "path",
Description: "object name and auth scope, such as for teams and projects",
Example: "default",
Required: true,
Schema: spec.StringProperty().UniqueValues(),
},
},
}
query.Post.Tags = []string{"DataSource"}
query.Post.Description = options.QueryDescription
query.Post.Parameters = nil //
query.Post.RequestBody = &spec3.RequestBody{
RequestBodyProps: spec3.RequestBodyProps{
Content: map[string]*spec3.MediaType{
+16
View File
@@ -50,6 +50,7 @@ type QueryAPIBuilder struct {
converter *expr.ResultConverter
queryTypes *query.QueryTypeDefinitionList
legacyDatasourceLookup service.LegacyDataSourceLookup
connections DataSourceConnectionProvider
}
func NewQueryAPIBuilder(
@@ -60,6 +61,7 @@ func NewQueryAPIBuilder(
registerer prometheus.Registerer,
tracer tracing.Tracer,
legacyDatasourceLookup service.LegacyDataSourceLookup,
connections DataSourceConnectionProvider,
) (*QueryAPIBuilder, error) {
// Include well typed query definitions
var queryTypes *query.QueryTypeDefinitionList
@@ -86,6 +88,7 @@ func NewQueryAPIBuilder(
tracer: tracer,
features: features,
queryTypes: queryTypes,
connections: connections,
converter: &expr.ResultConverter{
Features: features,
Tracer: tracer,
@@ -127,6 +130,8 @@ func RegisterAPIService(
return authorizer.DecisionAllow, "", nil
})
reg := client.NewDataSourceRegistryFromStore(pluginStore, dataSourcesService)
builder, err := NewQueryAPIBuilder(
features,
client.NewSingleTenantInstanceProvider(cfg, features, pluginClient, pCtxProvider, accessControl),
@@ -135,6 +140,7 @@ func RegisterAPIService(
registerer,
tracer,
legacyDatasourceLookup,
&connectionsProvider{dsService: dataSourcesService, registry: reg},
)
apiregistration.RegisterAPI(builder)
return builder, err
@@ -148,6 +154,8 @@ func addKnownTypes(scheme *runtime.Scheme, gv schema.GroupVersion) {
scheme.AddKnownTypes(gv,
&query.DataSourceApiServer{},
&query.DataSourceApiServerList{},
&query.DataSourceConnection{},
&query.DataSourceConnectionList{},
&query.QueryDataRequest{},
&query.QueryDataResponse{},
&query.QueryTypeDefinition{},
@@ -170,6 +178,14 @@ func (b *QueryAPIBuilder) UpdateAPIGroupInfo(apiGroupInfo *genericapiserver.APIG
storage := map[string]rest.Storage{}
// Get a list of all datasource instances
if b.features.IsEnabledGlobally(featuremgmt.FlagQueryServiceWithConnections) {
// Eventually this would be backed either by search or reconciler pattern
storage[query.ConnectionResourceInfo.StoragePath()] = &connectionAccess{
connections: b.connections,
}
}
plugins := newPluginsStorage(b.registry)
storage[plugins.resourceInfo.StoragePath()] = plugins
if !b.features.IsEnabledGlobally(featuremgmt.FlagGrafanaAPIServerWithExperimentalAPIs) {