From 5499ad8023519d492ca0ec056db6a3d45447e47b Mon Sep 17 00:00:00 2001 From: Dafydd Date: Thu, 4 Dec 2025 16:34:22 +0000 Subject: [PATCH] provide an interface for the datasourceConnection --- contribute/backend/style-guide.md | 2 + .../apis/collections/datasources_validator.go | 42 ++---- .../collections/datasources_validator_test.go | 80 +++++++++- pkg/registry/apis/collections/register.go | 6 +- pkg/server/wire.go | 2 + pkg/server/wire_gen.go | 9 +- .../datasources/service/client/client.go | 142 +++++++++--------- .../datasources/service/client/client_mock.go | 98 ++++++++++++ 8 files changed, 265 insertions(+), 116 deletions(-) create mode 100644 pkg/services/datasources/service/client/client_mock.go diff --git a/contribute/backend/style-guide.md b/contribute/backend/style-guide.md index dec2ea1db10..b7ab568451f 100644 --- a/contribute/backend/style-guide.md +++ b/contribute/backend/style-guide.md @@ -181,6 +181,8 @@ import ( //go:generate mockery --name InterfaceName --structname MockImplementationName --inpackage --filename my_implementation_mock.go ``` +The current `go:generate` command format used in this repository is only compatible with mockery v2. + ## Globals As a general rule of thumb, avoid using global variables, since they make the code difficult to maintain and reason diff --git a/pkg/registry/apis/collections/datasources_validator.go b/pkg/registry/apis/collections/datasources_validator.go index 7e8eab7d255..79226656caf 100644 --- a/pkg/registry/apis/collections/datasources_validator.go +++ b/pkg/registry/apis/collections/datasources_validator.go @@ -3,23 +3,21 @@ package collections import ( "context" "fmt" - "net/http" collections "github.com/grafana/grafana/apps/collections/pkg/apis/collections/v1alpha1" - "github.com/grafana/grafana/pkg/services/apiserver" "github.com/grafana/grafana/pkg/services/apiserver/builder" + "github.com/grafana/grafana/pkg/services/datasources/service/client" "k8s.io/apiserver/pkg/admission" - "k8s.io/client-go/kubernetes" ) var _ builder.APIGroupValidation = (*DatasourceStacksValidator)(nil) type DatasourceStacksValidator struct { - restConfigProvider apiserver.RestConfigProvider + dsClient client.DataSourceConnectionClient } -func GetDatasourceStacksValidator(restConfigProvider apiserver.RestConfigProvider) builder.APIGroupValidation { - return &DatasourceStacksValidator{restConfigProvider: restConfigProvider} +func GetDatasourceStacksValidator(dsClient client.DataSourceConnectionClient) builder.APIGroupValidation { + return &DatasourceStacksValidator{dsClient: dsClient} } func (v *DatasourceStacksValidator) Validate(ctx context.Context, a admission.Attributes, o admission.ObjectInterfaces) (err error) { @@ -66,13 +64,9 @@ func (v *DatasourceStacksValidator) Validate(ctx context.Context, a admission.At } exists, err := v.checkDatasourceExists(ctx, template[key].Group, item.DataSourceRef) - if err != nil { - return fmt.Errorf("error fetching: datasource '%s' does not exist (%s %s): %w", item.DataSourceRef, a.GetName(), a.GetKind().GroupVersion().String(), err) + if err != nil || !exists { + return fmt.Errorf("datasource '%s' in group '%s' does not exist (%s %s): %w", item.DataSourceRef, template[key].Group, a.GetName(), a.GetKind().GroupVersion().String(), err) } - if !exists { - return fmt.Errorf("datasource '%s' does not exist (%s %s)", item.DataSourceRef, a.GetName(), a.GetKind().GroupVersion().String()) - } - } } @@ -80,33 +74,15 @@ func (v *DatasourceStacksValidator) Validate(ctx context.Context, a admission.At } func (v *DatasourceStacksValidator) checkDatasourceExists(ctx context.Context, group, name string) (bool, error) { - cfg, err := v.restConfigProvider.GetRestConfig(ctx) + dsConn, err := v.dsClient.Get(ctx, group, "", name) if err != nil { return false, err } - client, err := kubernetes.NewForConfig(cfg) - if err != nil { - return false, err - } - - result := client.RESTClient().Get(). - Prefix("apis", group, "v0alpha1"). - Namespace("default"). - Resource("datasources"). - Name(name). - Do(ctx) - - if err = result.Error(); err != nil { - return false, err - } - - var statusCode int - - result = result.StatusCode(&statusCode) - if statusCode == http.StatusNotFound { + if dsConn == nil { return false, nil } return true, nil + } diff --git a/pkg/registry/apis/collections/datasources_validator_test.go b/pkg/registry/apis/collections/datasources_validator_test.go index 77b3d1644cd..2f154e9a330 100644 --- a/pkg/registry/apis/collections/datasources_validator_test.go +++ b/pkg/registry/apis/collections/datasources_validator_test.go @@ -5,23 +5,29 @@ import ( "testing" collectionsv1alpha1 "github.com/grafana/grafana/apps/collections/pkg/apis/collections/v1alpha1" + queryv0alpha1 "github.com/grafana/grafana/pkg/apis/query/v0alpha1" "github.com/grafana/grafana/pkg/registry/apis/collections" + datasourcesclient "github.com/grafana/grafana/pkg/services/datasources/service/client" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" "k8s.io/apimachinery/pkg/runtime" "k8s.io/apimachinery/pkg/runtime/schema" "k8s.io/apiserver/pkg/admission" ) func TestDataSourceValidator_Validate(t *testing.T) { - validator := &collections.DatasourceStacksValidator{} ctx := context.Background() tests := []struct { - name string - operation admission.Operation - object runtime.Object - expectError bool - errorMsg string + name string + operation admission.Operation + object runtime.Object + needMockDSClient bool // only set to true if you expect to make a call to the datasource client + dsClientReturnValue *queryv0alpha1.DataSourceConnection + dsClientReturnError error + expectError bool + errorMsg string }{ { name: "should return no error for invalid kind", @@ -94,6 +100,61 @@ func TestDataSourceValidator_Validate(t *testing.T) { expectError: true, errorMsg: "key 'notintemplate' is not in the DataSourceStack template (test-datasourcestack collections.grafana.app/v1alpha1)", }, + { + name: "error if data source does not exist", + operation: admission.Create, + object: &collectionsv1alpha1.DataSourceStack{ + Spec: collectionsv1alpha1.DataSourceStackSpec{ + Template: collectionsv1alpha1.DataSourceStackTemplateSpec{ + "key1": collectionsv1alpha1.DataSourceStackDataSourceStackTemplateItem{ + Name: "foo", + Group: "foo.grafana", + }, + }, + Modes: []collectionsv1alpha1.DataSourceStackModeSpec{ + { + Name: "prod", + Definition: collectionsv1alpha1.DataSourceStackMode{ + "key1": collectionsv1alpha1.DataSourceStackModeItem{ + DataSourceRef: "ref", + }, + }, + }, + }, + }, + }, + needMockDSClient: true, + dsClientReturnValue: nil, // no result - this is the default anyway + expectError: true, + errorMsg: "datasource 'ref' in group 'foo.grafana' does not exist (test-datasourcestack collections.grafana.app/v1alpha1)", + }, + { + name: "valid request", + operation: admission.Create, + object: &collectionsv1alpha1.DataSourceStack{ + Spec: collectionsv1alpha1.DataSourceStackSpec{ + Template: collectionsv1alpha1.DataSourceStackTemplateSpec{ + "key1": collectionsv1alpha1.DataSourceStackDataSourceStackTemplateItem{ + Name: "foo", + Group: "foo.grafana", + }, + }, + Modes: []collectionsv1alpha1.DataSourceStackModeSpec{ + { + Name: "prod", + Definition: collectionsv1alpha1.DataSourceStackMode{ + "key1": collectionsv1alpha1.DataSourceStackModeItem{ + DataSourceRef: "ref", + }, + }, + }, + }, + }, + }, + needMockDSClient: true, + dsClientReturnValue: &queryv0alpha1.DataSourceConnection{}, // returning any non-nil value will pass validation + expectError: false, + }, } for _, tt := range tests { @@ -105,6 +166,13 @@ func TestDataSourceValidator_Validate(t *testing.T) { Kind: schema.GroupVersionKind{Group: "collections.grafana.app", Version: "v1alpha1", Kind: "DataSourceStack"}, } + var client *datasourcesclient.MockDataSourceConnectionClient + if tt.needMockDSClient { + client = datasourcesclient.NewMockDataSourceConnectionClient(t) + client.On("Get", mock.Anything, mock.Anything, mock.Anything, mock.Anything).Return(tt.dsClientReturnValue, tt.dsClientReturnError) + } + + validator := collections.GetDatasourceStacksValidator(client) err := validator.Validate(ctx, attrs, nil) if tt.expectError { diff --git a/pkg/registry/apis/collections/register.go b/pkg/registry/apis/collections/register.go index a53481d826c..385b4014385 100644 --- a/pkg/registry/apis/collections/register.go +++ b/pkg/registry/apis/collections/register.go @@ -25,6 +25,7 @@ import ( "github.com/grafana/grafana/pkg/services/apiserver" "github.com/grafana/grafana/pkg/services/apiserver/builder" "github.com/grafana/grafana/pkg/services/apiserver/endpoints/request" + datasourcesClient "github.com/grafana/grafana/pkg/services/datasources/service/client" "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/services/star" "github.com/grafana/grafana/pkg/services/user" @@ -51,6 +52,7 @@ func RegisterAPIService( stars star.Service, users user.Service, apiregistration builder.APIRegistrar, + dsConnClientFactory datasourcesClient.DataSourceConnectionClientFactory, restConfigProvider apiserver.RestConfigProvider, ) *APIBuilder { // Requires development settings and clearly experimental @@ -59,9 +61,11 @@ func RegisterAPIService( return nil } + dsConnClient := dsConnClientFactory(restConfigProvider) + sql := legacy.NewLegacySQL(legacysql.NewDatabaseProvider(db)) builder := &APIBuilder{ - datasourceStacksValidator: GetDatasourceStacksValidator(restConfigProvider), + datasourceStacksValidator: GetDatasourceStacksValidator(dsConnClient), authorizer: &utils.AuthorizeFromName{ Resource: map[string][]utils.ResourceOwner{ "stars": {utils.UserResourceOwner}, diff --git a/pkg/server/wire.go b/pkg/server/wire.go index 0d4bb10c0b5..e3be194df68 100644 --- a/pkg/server/wire.go +++ b/pkg/server/wire.go @@ -88,6 +88,7 @@ import ( "github.com/grafana/grafana/pkg/services/datasourceproxy" "github.com/grafana/grafana/pkg/services/datasources" datasourceservice "github.com/grafana/grafana/pkg/services/datasources/service" + datasourcesclient "github.com/grafana/grafana/pkg/services/datasources/service/client" "github.com/grafana/grafana/pkg/services/dsquerierclient" "github.com/grafana/grafana/pkg/services/encryption" encryptionservice "github.com/grafana/grafana/pkg/services/encryption/service" @@ -476,6 +477,7 @@ var wireBasicSet = wire.NewSet( appregistry.WireSet, // Dashboard Kubernetes helpers dashboardclient.ProvideK8sClientWithFallback, + datasourcesclient.ProvideDataSourceConnectionClientFactory, ) var wireSet = wire.NewSet( diff --git a/pkg/server/wire_gen.go b/pkg/server/wire_gen.go index 9baae3adde6..d34b66c1527 100644 --- a/pkg/server/wire_gen.go +++ b/pkg/server/wire_gen.go @@ -134,6 +134,7 @@ import ( "github.com/grafana/grafana/pkg/services/datasources" "github.com/grafana/grafana/pkg/services/datasources/guardian" service9 "github.com/grafana/grafana/pkg/services/datasources/service" + client2 "github.com/grafana/grafana/pkg/services/datasources/service/client" "github.com/grafana/grafana/pkg/services/dsquerierclient" encryption2 "github.com/grafana/grafana/pkg/services/encryption" "github.com/grafana/grafana/pkg/services/encryption/provider" @@ -886,7 +887,8 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api } userStorageAPIBuilder := userstorage.RegisterAPIService(featureToggles, apiserverService, registerer) apiBuilder := preferences.RegisterAPIService(cfg, featureToggles, sqlStore, prefService, userService, apiserverService) - collectionsAPIBuilder := collections.RegisterAPIService(cfg, featureToggles, sqlStore, starService, userService, apiserverService, eventualRestConfigProvider) + dataSourceConnectionClientFactory := client2.ProvideDataSourceConnectionClientFactory(eventualRestConfigProvider) + collectionsAPIBuilder := collections.RegisterAPIService(cfg, featureToggles, sqlStore, starService, userService, apiserverService, dataSourceConnectionClientFactory, eventualRestConfigProvider) webhookExtraBuilder := webhooks.ProvideWebhooksWithImages(cfg, renderingService, resourceClient, eventualRestConfigProvider, registerer) v3 := extras.ProvideProvisioningExtraAPIs(webhookExtraBuilder) pullRequestWorker := pullrequest.ProvidePullRequestWorker(cfg, renderingService, resourceClient, eventualRestConfigProvider, registerer) @@ -1540,7 +1542,8 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac } userStorageAPIBuilder := userstorage.RegisterAPIService(featureToggles, apiserverService, registerer) apiBuilder := preferences.RegisterAPIService(cfg, featureToggles, sqlStore, prefService, userService, apiserverService) - collectionsAPIBuilder := collections.RegisterAPIService(cfg, featureToggles, sqlStore, starService, userService, apiserverService, eventualRestConfigProvider) + dataSourceConnectionClientFactory := client2.ProvideDataSourceConnectionClientFactory(eventualRestConfigProvider) + collectionsAPIBuilder := collections.RegisterAPIService(cfg, featureToggles, sqlStore, starService, userService, apiserverService, dataSourceConnectionClientFactory, eventualRestConfigProvider) webhookExtraBuilder := webhooks.ProvideWebhooksWithImages(cfg, renderingService, resourceClient, eventualRestConfigProvider, registerer) v3 := extras.ProvideProvisioningExtraAPIs(webhookExtraBuilder) pullRequestWorker := pullrequest.ProvidePullRequestWorker(cfg, renderingService, resourceClient, eventualRestConfigProvider, registerer) @@ -1785,7 +1788,7 @@ var withOTelSet = wire.NewSet( otelTracer, grpcserver.ProvideService, interceptors.ProvideAuthenticator, ) -var wireBasicSet = wire.NewSet(annotationsimpl.ProvideService, wire.Bind(new(annotations.Repository), new(*annotationsimpl.RepositoryImpl)), New, api.ProvideHTTPServer, query.ProvideService, wire.Bind(new(query.Service), new(*query.ServiceImpl)), bus.ProvideBus, wire.Bind(new(bus.Bus), new(*bus.InProcBus)), rendering.ProvideService, wire.Bind(new(rendering.Service), new(*rendering.RenderingService)), routing.ProvideRegister, wire.Bind(new(routing.RouteRegister), new(*routing.RouteRegisterImpl)), hooks.ProvideService, kvstore.ProvideService, localcache.ProvideService, bundleregistry.ProvideService, wire.Bind(new(supportbundles.Service), new(*bundleregistry.Service)), updatemanager.ProvideGrafanaService, updatemanager.ProvidePluginsService, service.ProvideService, wire.Bind(new(usagestats.Service), new(*service.UsageStats)), validator3.ProvideService, provisioning.ProvideStubProvisioningService, legacy.ProvideMigratorDashboardAccessor, migrations2.ProvideUnifiedMigrator, pluginsintegration.WireSet, dashboards.ProvideFileStoreManager, wire.Bind(new(dashboards.FileStore), new(*dashboards.FileStoreManager)), cloudwatch.ProvideService, cloudmonitoring.ProvideService, azuremonitor.ProvideService, postgres.ProvideService, mysql.ProvideService, mssql.ProvideService, store.ProvideEntityEventsService, dualwrite.ProvideService, httpclientprovider.New, wire.Bind(new(httpclient.Provider), new(*httpclient2.Provider)), serverlock.ProvideService, wire.Bind(new(installsync.ServerLock), new(*serverlock.ServerLockService)), annotationsimpl.ProvideCleanupService, wire.Bind(new(annotations.Cleaner), new(*annotationsimpl.CleanupServiceImpl)), cleanup.ProvideService, shorturlimpl.ProvideService, wire.Bind(new(shorturls.Service), new(*shorturlimpl.ShortURLService)), queryhistory.ProvideService, wire.Bind(new(queryhistory.Service), new(*queryhistory.QueryHistoryService)), correlations.ProvideService, wire.Bind(new(correlations.Service), new(*correlations.CorrelationsService)), quotaimpl.ProvideService, remotecache.ProvideService, wire.Bind(new(remotecache.CacheStorage), new(*remotecache.RemoteCache)), authinfoimpl.ProvideService, wire.Bind(new(login.AuthInfoService), new(*authinfoimpl.Service)), authinfoimpl.ProvideStore, datasourceproxy.ProvideService, sort.ProvideService, search2.ProvideService, searchV2.ProvideService, searchV2.ProvideSearchHTTPService, store.ProvideService, store.ProvideSystemUsersService, live.ProvideService, pushhttp.ProvideService, contexthandler.ProvideService, service12.ProvideService, wire.Bind(new(service12.LDAP), new(*service12.LDAPImpl)), jwt.ProvideService, wire.Bind(new(jwt.JWTService), new(*jwt.AuthService)), store2.ProvideDBStore, image.ProvideDeleteExpiredService, ngalert.ProvideService, librarypanels.ProvideService, wire.Bind(new(librarypanels.Service), new(*librarypanels.LibraryPanelService)), libraryelements.ProvideService, wire.Bind(new(libraryelements.Service), new(*libraryelements.LibraryElementService)), notifications.ProvideService, notifications.ProvideSmtpService, github.ProvideFactory, tracing.ProvideService, tracing.ProvideTracingConfig, wire.Bind(new(tracing.Tracer), new(*tracing.TracingService)), withOTelSet, testdatasource.ProvideService, api4.ProvideService, opentsdb.ProvideService, socialimpl.ProvideService, influxdb.ProvideService, wire.Bind(new(social.Service), new(*socialimpl.SocialService)), tempo.ProvideService, loki.ProvideService, graphite.ProvideService, prometheus.ProvideService, elasticsearch.ProvideService, pyroscope.ProvideService, parca.ProvideService, zipkin.ProvideService, jaeger.ProvideService, service9.ProvideCacheService, wire.Bind(new(datasources.CacheService), new(*service9.CacheServiceImpl)), service2.ProvideEncryptionService, wire.Bind(new(encryption2.Internal), new(*service2.Service)), manager.ProvideSecretsService, wire.Bind(new(secrets.Service), new(*manager.SecretsService)), database.ProvideSecretsStore, wire.Bind(new(secrets.Store), new(*database.SecretsStoreImpl)), garbagecollectionworker.ProvideWorker, grafanads.ProvideService, wire.Bind(new(dashboardsnapshots.Store), new(*database5.DashboardSnapshotStore)), database5.ProvideStore, wire.Bind(new(dashboardsnapshots.Service), new(*service10.ServiceImpl)), service10.ProvideService, service9.ProvideService, wire.Bind(new(datasources.DataSourceService), new(*service9.Service)), service9.ProvideLegacyDataSourceLookup, retriever.ProvideService, wire.Bind(new(serviceaccounts.ServiceAccountRetriever), new(*retriever.Service)), ossaccesscontrol.ProvideServiceAccountPermissions, wire.Bind(new(accesscontrol.ServiceAccountPermissionsService), new(*ossaccesscontrol.ServiceAccountPermissionsService)), manager3.ProvideServiceAccountsService, proxy.ProvideServiceAccountsProxy, wire.Bind(new(serviceaccounts.Service), new(*proxy.ServiceAccountsProxy)), dsquerierclient.NewNullQSDatasourceClientBuilder, expr.ProvideService, featuremgmt.ProvideManagerService, featuremgmt.ProvideToggles, service7.ProvideDashboardServiceImpl, wire.Bind(new(dashboards2.PermissionsRegistrationService), new(*service7.DashboardServiceImpl)), service7.ProvideDashboardService, service7.ProvideDashboardProvisioningService, service7.ProvideDashboardPluginService, database2.ProvideDashboardStore, folderimpl.ProvideService, wire.Bind(new(folder.Service), new(*folderimpl.Service)), wire.Bind(new(folder.LegacyService), new(*folderimpl.Service)), folderimpl.ProvideStore, wire.Bind(new(folder.Store), new(*folderimpl.FolderStoreImpl)), service11.ProvideService, wire.Bind(new(dashboardimport.Service), new(*service11.ImportDashboardService)), service8.ProvideService, wire.Bind(new(plugindashboards.Service), new(*service8.Service)), service8.ProvideDashboardUpdater, kvstore2.ProvideService, avatar.ProvideAvatarCacheServer, statscollector.ProvideService, csrf.ProvideCSRFFilter, wire.Bind(new(csrf.Service), new(*csrf.CSRF)), ossaccesscontrol.ProvideTeamPermissions, wire.Bind(new(accesscontrol.TeamPermissionsService), new(*ossaccesscontrol.TeamPermissionsService)), ossaccesscontrol.ProvideFolderPermissions, wire.Bind(new(accesscontrol.FolderPermissionsService), new(*ossaccesscontrol.FolderPermissionsService)), ossaccesscontrol.ProvideDashboardPermissions, wire.Bind(new(accesscontrol.DashboardPermissionsService), new(*ossaccesscontrol.DashboardPermissionsService)), ossaccesscontrol.ProvideReceiverPermissionsService, wire.Bind(new(accesscontrol.ReceiverPermissionsService), new(*ossaccesscontrol.ReceiverPermissionsService)), starimpl.ProvideService, playlistimpl.ProvideService, apikeyimpl.ProvideService, dashverimpl.ProvideService, service3.ProvideService, wire.Bind(new(publicdashboards.Service), new(*service3.PublicDashboardServiceImpl)), database3.ProvideStore, wire.Bind(new(publicdashboards.Store), new(*database3.PublicDashboardStoreImpl)), metric.ProvideService, api2.ProvideApi, api3.ProvideApi, userimpl.ProvideService, orgimpl.ProvideService, orgimpl.ProvideDeletionService, statsimpl.ProvideService, grpccontext.ProvideContextHandler, grpcserver.ProvideHealthService, grpcserver.ProvideReflectionService, resolver.ProvideEntityReferenceResolver, teamimpl.ProvideService, teamapi.ProvideTeamAPI, tempuserimpl.ProvideService, loginattemptimpl.ProvideService, wire.Bind(new(loginattempt.Service), new(*loginattemptimpl.Service)), migrations3.ProvideDataSourceMigrationService, migrations3.ProvideSecretMigrationProvider, wire.Bind(new(migrations3.SecretMigrationProvider), new(*migrations3.SecretMigrationProviderImpl)), promtypemigration.ProvideAzurePromMigrationService, promtypemigration.ProvideAmazonPromMigrationService, promtypemigration.ProvidePromTypeMigrationProvider, wire.Bind(new(promtypemigration.PromTypeMigrationProvider), new(*promtypemigration.PromTypeMigrationProviderImpl)), resourcepermissions.NewActionSetService, wire.Bind(new(accesscontrol.ActionResolver), new(resourcepermissions.ActionSetService)), wire.Bind(new(pluginaccesscontrol.ActionSetRegistry), new(resourcepermissions.ActionSetService)), permreg.ProvidePermissionRegistry, acimpl.ProvideAccessControl, accesscontrol.ProvideFixedRolesLoader, dualwrite2.ProvideZanzanaReconciler, navtreeimpl.ProvideService, wire.Bind(new(accesscontrol.AccessControl), new(*acimpl.AccessControl)), wire.Bind(new(notifications.TempUserStore), new(tempuser.Service)), tagimpl.ProvideService, wire.Bind(new(tag.Service), new(*tagimpl.Service)), authnimpl.ProvideService, authnimpl.ProvideIdentitySynchronizer, authnimpl.ProvideAuthnService, authnimpl.ProvideAuthnServiceAuthenticateOnly, authnimpl.ProvideRegistration, supportbundlesimpl.ProvideService, extsvcaccounts.ProvideExtSvcAccountsService, wire.Bind(new(serviceaccounts.ExtSvcAccountsService), new(*extsvcaccounts.ExtSvcAccountsService)), registry2.ProvideExtSvcRegistry, wire.Bind(new(extsvcauth.ExternalServiceRegistry), new(*registry2.Registry)), anonstore.ProvideAnonDBStore, wire.Bind(new(anonstore.AnonStore), new(*anonstore.AnonDBStore)), loggermw.Provide, slogadapter.Provide, signingkeysimpl.ProvideEmbeddedSigningKeysService, wire.Bind(new(signingkeys.Service), new(*signingkeysimpl.Service)), ssosettingsimpl.ProvideService, wire.Bind(new(ssosettings.Service), new(*ssosettingsimpl.Service)), idimpl.ProvideService, wire.Bind(new(auth.IDService), new(*idimpl.Service)), cloudmigrationimpl.ProvideService, caching.ProvideCachingServiceClient, userimpl.ProvideVerifier, connectors.ProvideOrgRoleMapper, wire.Bind(new(user.Verifier), new(*userimpl.Verifier)), authz.WireSet, metadata.ProvideSecureValueMetadataStorage, metadata.ProvideKeeperMetadataStorage, metadata.ProvideDecryptStorage, decrypt.ProvideDecryptAuthorizer, wire.Value([]decrypt.ExtraOwnerDecrypter(nil)), decrypt.ProvideDecryptService, inline.ProvideInlineSecureValueService, encryption.ProvideDataKeyStorage, encryption.ProvideGlobalDataKeyStorage, encryption.ProvideEncryptedValueStorage, encryption.ProvideGlobalEncryptedValueStorage, encryption.ProvideEncryptedValueMigrationExecutor, service5.ProvideSecureValueService, validator.ProvideKeeperValidator, validator.ProvideSecureValueValidator, mutator.ProvideKeeperMutator, mutator.ProvideSecureValueMutator, migrator.NewWithEngine, database4.ProvideDatabase, clock.ProvideClock, wire.Bind(new(contracts.Database), new(*database4.Database)), wire.Bind(new(contracts.Clock), new(*clock.Clock)), manager2.ProvideEncryptionManager, service4.ProvideAESGCMCipherService, resource.ProvideStorageMetrics, resource.ProvideIndexMetrics, migrations2.ProvideUnifiedStorageMigrationService, apiserver.WireSet, apiregistry.WireSet, appregistry.WireSet, client.ProvideK8sClientWithFallback) +var wireBasicSet = wire.NewSet(annotationsimpl.ProvideService, wire.Bind(new(annotations.Repository), new(*annotationsimpl.RepositoryImpl)), New, api.ProvideHTTPServer, query.ProvideService, wire.Bind(new(query.Service), new(*query.ServiceImpl)), bus.ProvideBus, wire.Bind(new(bus.Bus), new(*bus.InProcBus)), rendering.ProvideService, wire.Bind(new(rendering.Service), new(*rendering.RenderingService)), routing.ProvideRegister, wire.Bind(new(routing.RouteRegister), new(*routing.RouteRegisterImpl)), hooks.ProvideService, kvstore.ProvideService, localcache.ProvideService, bundleregistry.ProvideService, wire.Bind(new(supportbundles.Service), new(*bundleregistry.Service)), updatemanager.ProvideGrafanaService, updatemanager.ProvidePluginsService, service.ProvideService, wire.Bind(new(usagestats.Service), new(*service.UsageStats)), validator3.ProvideService, provisioning.ProvideStubProvisioningService, legacy.ProvideMigratorDashboardAccessor, migrations2.ProvideUnifiedMigrator, pluginsintegration.WireSet, dashboards.ProvideFileStoreManager, wire.Bind(new(dashboards.FileStore), new(*dashboards.FileStoreManager)), cloudwatch.ProvideService, cloudmonitoring.ProvideService, azuremonitor.ProvideService, postgres.ProvideService, mysql.ProvideService, mssql.ProvideService, store.ProvideEntityEventsService, dualwrite.ProvideService, httpclientprovider.New, wire.Bind(new(httpclient.Provider), new(*httpclient2.Provider)), serverlock.ProvideService, wire.Bind(new(installsync.ServerLock), new(*serverlock.ServerLockService)), annotationsimpl.ProvideCleanupService, wire.Bind(new(annotations.Cleaner), new(*annotationsimpl.CleanupServiceImpl)), cleanup.ProvideService, shorturlimpl.ProvideService, wire.Bind(new(shorturls.Service), new(*shorturlimpl.ShortURLService)), queryhistory.ProvideService, wire.Bind(new(queryhistory.Service), new(*queryhistory.QueryHistoryService)), correlations.ProvideService, wire.Bind(new(correlations.Service), new(*correlations.CorrelationsService)), quotaimpl.ProvideService, remotecache.ProvideService, wire.Bind(new(remotecache.CacheStorage), new(*remotecache.RemoteCache)), authinfoimpl.ProvideService, wire.Bind(new(login.AuthInfoService), new(*authinfoimpl.Service)), authinfoimpl.ProvideStore, datasourceproxy.ProvideService, sort.ProvideService, search2.ProvideService, searchV2.ProvideService, searchV2.ProvideSearchHTTPService, store.ProvideService, store.ProvideSystemUsersService, live.ProvideService, pushhttp.ProvideService, contexthandler.ProvideService, service12.ProvideService, wire.Bind(new(service12.LDAP), new(*service12.LDAPImpl)), jwt.ProvideService, wire.Bind(new(jwt.JWTService), new(*jwt.AuthService)), store2.ProvideDBStore, image.ProvideDeleteExpiredService, ngalert.ProvideService, librarypanels.ProvideService, wire.Bind(new(librarypanels.Service), new(*librarypanels.LibraryPanelService)), libraryelements.ProvideService, wire.Bind(new(libraryelements.Service), new(*libraryelements.LibraryElementService)), notifications.ProvideService, notifications.ProvideSmtpService, github.ProvideFactory, tracing.ProvideService, tracing.ProvideTracingConfig, wire.Bind(new(tracing.Tracer), new(*tracing.TracingService)), withOTelSet, testdatasource.ProvideService, api4.ProvideService, opentsdb.ProvideService, socialimpl.ProvideService, influxdb.ProvideService, wire.Bind(new(social.Service), new(*socialimpl.SocialService)), tempo.ProvideService, loki.ProvideService, graphite.ProvideService, prometheus.ProvideService, elasticsearch.ProvideService, pyroscope.ProvideService, parca.ProvideService, zipkin.ProvideService, jaeger.ProvideService, service9.ProvideCacheService, wire.Bind(new(datasources.CacheService), new(*service9.CacheServiceImpl)), service2.ProvideEncryptionService, wire.Bind(new(encryption2.Internal), new(*service2.Service)), manager.ProvideSecretsService, wire.Bind(new(secrets.Service), new(*manager.SecretsService)), database.ProvideSecretsStore, wire.Bind(new(secrets.Store), new(*database.SecretsStoreImpl)), garbagecollectionworker.ProvideWorker, grafanads.ProvideService, wire.Bind(new(dashboardsnapshots.Store), new(*database5.DashboardSnapshotStore)), database5.ProvideStore, wire.Bind(new(dashboardsnapshots.Service), new(*service10.ServiceImpl)), service10.ProvideService, service9.ProvideService, wire.Bind(new(datasources.DataSourceService), new(*service9.Service)), service9.ProvideLegacyDataSourceLookup, retriever.ProvideService, wire.Bind(new(serviceaccounts.ServiceAccountRetriever), new(*retriever.Service)), ossaccesscontrol.ProvideServiceAccountPermissions, wire.Bind(new(accesscontrol.ServiceAccountPermissionsService), new(*ossaccesscontrol.ServiceAccountPermissionsService)), manager3.ProvideServiceAccountsService, proxy.ProvideServiceAccountsProxy, wire.Bind(new(serviceaccounts.Service), new(*proxy.ServiceAccountsProxy)), dsquerierclient.NewNullQSDatasourceClientBuilder, expr.ProvideService, featuremgmt.ProvideManagerService, featuremgmt.ProvideToggles, service7.ProvideDashboardServiceImpl, wire.Bind(new(dashboards2.PermissionsRegistrationService), new(*service7.DashboardServiceImpl)), service7.ProvideDashboardService, service7.ProvideDashboardProvisioningService, service7.ProvideDashboardPluginService, database2.ProvideDashboardStore, folderimpl.ProvideService, wire.Bind(new(folder.Service), new(*folderimpl.Service)), wire.Bind(new(folder.LegacyService), new(*folderimpl.Service)), folderimpl.ProvideStore, wire.Bind(new(folder.Store), new(*folderimpl.FolderStoreImpl)), service11.ProvideService, wire.Bind(new(dashboardimport.Service), new(*service11.ImportDashboardService)), service8.ProvideService, wire.Bind(new(plugindashboards.Service), new(*service8.Service)), service8.ProvideDashboardUpdater, kvstore2.ProvideService, avatar.ProvideAvatarCacheServer, statscollector.ProvideService, csrf.ProvideCSRFFilter, wire.Bind(new(csrf.Service), new(*csrf.CSRF)), ossaccesscontrol.ProvideTeamPermissions, wire.Bind(new(accesscontrol.TeamPermissionsService), new(*ossaccesscontrol.TeamPermissionsService)), ossaccesscontrol.ProvideFolderPermissions, wire.Bind(new(accesscontrol.FolderPermissionsService), new(*ossaccesscontrol.FolderPermissionsService)), ossaccesscontrol.ProvideDashboardPermissions, wire.Bind(new(accesscontrol.DashboardPermissionsService), new(*ossaccesscontrol.DashboardPermissionsService)), ossaccesscontrol.ProvideReceiverPermissionsService, wire.Bind(new(accesscontrol.ReceiverPermissionsService), new(*ossaccesscontrol.ReceiverPermissionsService)), starimpl.ProvideService, playlistimpl.ProvideService, apikeyimpl.ProvideService, dashverimpl.ProvideService, service3.ProvideService, wire.Bind(new(publicdashboards.Service), new(*service3.PublicDashboardServiceImpl)), database3.ProvideStore, wire.Bind(new(publicdashboards.Store), new(*database3.PublicDashboardStoreImpl)), metric.ProvideService, api2.ProvideApi, api3.ProvideApi, userimpl.ProvideService, orgimpl.ProvideService, orgimpl.ProvideDeletionService, statsimpl.ProvideService, grpccontext.ProvideContextHandler, grpcserver.ProvideHealthService, grpcserver.ProvideReflectionService, resolver.ProvideEntityReferenceResolver, teamimpl.ProvideService, teamapi.ProvideTeamAPI, tempuserimpl.ProvideService, loginattemptimpl.ProvideService, wire.Bind(new(loginattempt.Service), new(*loginattemptimpl.Service)), migrations3.ProvideDataSourceMigrationService, migrations3.ProvideSecretMigrationProvider, wire.Bind(new(migrations3.SecretMigrationProvider), new(*migrations3.SecretMigrationProviderImpl)), promtypemigration.ProvideAzurePromMigrationService, promtypemigration.ProvideAmazonPromMigrationService, promtypemigration.ProvidePromTypeMigrationProvider, wire.Bind(new(promtypemigration.PromTypeMigrationProvider), new(*promtypemigration.PromTypeMigrationProviderImpl)), resourcepermissions.NewActionSetService, wire.Bind(new(accesscontrol.ActionResolver), new(resourcepermissions.ActionSetService)), wire.Bind(new(pluginaccesscontrol.ActionSetRegistry), new(resourcepermissions.ActionSetService)), permreg.ProvidePermissionRegistry, acimpl.ProvideAccessControl, accesscontrol.ProvideFixedRolesLoader, dualwrite2.ProvideZanzanaReconciler, navtreeimpl.ProvideService, wire.Bind(new(accesscontrol.AccessControl), new(*acimpl.AccessControl)), wire.Bind(new(notifications.TempUserStore), new(tempuser.Service)), tagimpl.ProvideService, wire.Bind(new(tag.Service), new(*tagimpl.Service)), authnimpl.ProvideService, authnimpl.ProvideIdentitySynchronizer, authnimpl.ProvideAuthnService, authnimpl.ProvideAuthnServiceAuthenticateOnly, authnimpl.ProvideRegistration, supportbundlesimpl.ProvideService, extsvcaccounts.ProvideExtSvcAccountsService, wire.Bind(new(serviceaccounts.ExtSvcAccountsService), new(*extsvcaccounts.ExtSvcAccountsService)), registry2.ProvideExtSvcRegistry, wire.Bind(new(extsvcauth.ExternalServiceRegistry), new(*registry2.Registry)), anonstore.ProvideAnonDBStore, wire.Bind(new(anonstore.AnonStore), new(*anonstore.AnonDBStore)), loggermw.Provide, slogadapter.Provide, signingkeysimpl.ProvideEmbeddedSigningKeysService, wire.Bind(new(signingkeys.Service), new(*signingkeysimpl.Service)), ssosettingsimpl.ProvideService, wire.Bind(new(ssosettings.Service), new(*ssosettingsimpl.Service)), idimpl.ProvideService, wire.Bind(new(auth.IDService), new(*idimpl.Service)), cloudmigrationimpl.ProvideService, caching.ProvideCachingServiceClient, userimpl.ProvideVerifier, connectors.ProvideOrgRoleMapper, wire.Bind(new(user.Verifier), new(*userimpl.Verifier)), authz.WireSet, metadata.ProvideSecureValueMetadataStorage, metadata.ProvideKeeperMetadataStorage, metadata.ProvideDecryptStorage, decrypt.ProvideDecryptAuthorizer, wire.Value([]decrypt.ExtraOwnerDecrypter(nil)), decrypt.ProvideDecryptService, inline.ProvideInlineSecureValueService, encryption.ProvideDataKeyStorage, encryption.ProvideGlobalDataKeyStorage, encryption.ProvideEncryptedValueStorage, encryption.ProvideGlobalEncryptedValueStorage, encryption.ProvideEncryptedValueMigrationExecutor, service5.ProvideSecureValueService, validator.ProvideKeeperValidator, validator.ProvideSecureValueValidator, mutator.ProvideKeeperMutator, mutator.ProvideSecureValueMutator, migrator.NewWithEngine, database4.ProvideDatabase, clock.ProvideClock, wire.Bind(new(contracts.Database), new(*database4.Database)), wire.Bind(new(contracts.Clock), new(*clock.Clock)), manager2.ProvideEncryptionManager, service4.ProvideAESGCMCipherService, resource.ProvideStorageMetrics, resource.ProvideIndexMetrics, migrations2.ProvideUnifiedStorageMigrationService, apiserver.WireSet, apiregistry.WireSet, appregistry.WireSet, client.ProvideK8sClientWithFallback, client2.ProvideDataSourceConnectionClientFactory) var wireSet = wire.NewSet( wireBasicSet, metrics.WireSet, sqlstore.ProvideService, metrics2.ProvideService, wire.Bind(new(notifications.Service), new(*notifications.NotificationService)), wire.Bind(new(notifications.WebhookSender), new(*notifications.NotificationService)), wire.Bind(new(notifications.EmailSender), new(*notifications.NotificationService)), wire.Bind(new(db.DB), new(*sqlstore.SQLStore)), prefimpl.ProvideService, oauthtoken.ProvideService, wire.Bind(new(oauthtoken.OAuthTokenService), new(*oauthtoken.Service)), wire.Bind(new(cleanup.AlertRuleService), new(*store2.DBstore)), diff --git a/pkg/services/datasources/service/client/client.go b/pkg/services/datasources/service/client/client.go index 9e3d460b6e6..5a8a80374cb 100644 --- a/pkg/services/datasources/service/client/client.go +++ b/pkg/services/datasources/service/client/client.go @@ -1,90 +1,86 @@ package client -// import ( -// "context" -// "encoding/json" +import ( + "context" + "errors" + "net/http" -// "github.com/grafana/grafana/pkg/services/apiserver" -// "github.com/grafana/grafana/pkg/services/apiserver/client" -// v1 "k8s.io/apimachinery/pkg/apis/meta/v1" -// "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" -// "k8s.io/apimachinery/pkg/runtime/schema" -// "k8s.io/client-go/kubernetes" -// ) + datasourcev0alpha1 "github.com/grafana/grafana/pkg/apis/datasource/v0alpha1" + queryv0alpha1 "github.com/grafana/grafana/pkg/apis/query/v0alpha1" + "github.com/grafana/grafana/pkg/services/apiserver" + "k8s.io/client-go/kubernetes" +) -// // K8sClientFactory creates a K8sClient for the given group -// type K8sClientFactory func(ctx context.Context, group string, version string) client.K8sHandler +// DataSourceConnectionClient can get information about data source connections. +// +//go:generate mockery --name DataSourceConnectionClient --structname MockDataSourceConnectionClient --inpackage --filename=client_mock.go --with-expecter +type DataSourceConnectionClient interface { + Get(ctx context.Context, group, version, name string) (*queryv0alpha1.DataSourceConnection, error) +} -// type K8sHandler struct { -// client.K8sHandler -// restConfigProvider apiserver.RestConfigProvider -// gvr schema.GroupVersionResource -// } +func ProvideDataSourceConnectionClientFactory( + restConfigProvider apiserver.RestConfigProvider, +) DataSourceConnectionClientFactory { + return func(configProvider apiserver.RestConfigProvider) DataSourceConnectionClient { + return &dataSourceConnectionClient{ + configProvider: configProvider, + } + } +} -// type K8sClient struct { -// client.K8sHandler -// newClientFunc K8sClientFactory -// } +type DataSourceConnectionClientFactory func(configProvider apiserver.RestConfigProvider) DataSourceConnectionClient -// func ProvideK8sClient( -// restConfigProvider apiserver.RestConfigProvider, -// ) K8sHandler { -// return NewK8sClient(restConfigProvider) -// } +type dataSourceConnectionClient struct { + configProvider apiserver.RestConfigProvider +} -// func NewK8sClient(restConfigProvider apiserver.RestConfigProvider) *K8sClient { -// newClientFunc := newK8sClientFactory(restConfigProvider) -// return &K8sClient{ -// K8sHandler: newClientFunc(context.Background(),), -// } -// } +func (dc *dataSourceConnectionClient) Get(ctx context.Context, group, version, name string) (*queryv0alpha1.DataSourceConnection, error) { + cfg, err := dc.configProvider.GetRestConfig(ctx) + if err != nil { + return nil, err + } -// func (c K8sHandler) Get(ctx context.Context, name string, orgID int64, options v1.GetOptions, subresource ...string) (*unstructured.Unstructured, error) { -// cfg, err := c.restConfigProvider.GetRestConfig(ctx) -// if err != nil { -// return nil, err -// } -// client, err := kubernetes.NewForConfig(cfg) -// if err != nil { -// return nil, err -// } + client, err := kubernetes.NewForConfig(cfg) + if err != nil { + return nil, err + } -// result := client.RESTClient().Get(). -// Prefix("apis", c.gvr.Group, c.gvr.Version). -// Namespace(string(orgID)). -// Resource(c.gvr.Resource). -// Name(name). -// Do(ctx) + if version == "" { + version = "v0alpha1" + } -// if err = result.Error(); err != nil { -// return nil, err -// } + result := client.RESTClient().Get(). + Prefix("apis", group, version). + Namespace("default"). // TODO do something about namespace + Resource("datasources"). + Name(name). + Do(ctx) -// body, err := result.Raw() -// if err != nil { -// return nil, err -// } + if err = result.Error(); err != nil { + return nil, err + } -// value := &unstructured.Unstructured{} -// if err = json.Unmarshal(body, value); err != nil { -// return nil, err -// } + var statusCode int -// return value, nil -// } + result = result.StatusCode(&statusCode) + if statusCode == http.StatusNotFound { + return nil, errors.New("not found") + } -// func newK8sClientFactory(restConfigProvider apiserver.RestConfigProvider) K8sClientFactory { -// return func(ctx context.Context, group string, version string) client.K8sHandler { -// gvr := schema.GroupVersionResource{ -// Group: group, -// Version: version, -// Resource: "datasources", -// } + fullDS := datasourcev0alpha1.DataSource{} + err = result.Into(&fullDS) + if err != nil { + return nil, err + } -// return K8sHandler{ -// restConfigProvider: restConfigProvider, -// gvr: gvr, -// } + dsConnection := &queryv0alpha1.DataSourceConnection{ + Title: fullDS.Spec.Title(), + Datasource: queryv0alpha1.DataSourceConnectionRef{ + Group: fullDS.GroupVersionKind().Group, + Name: fullDS.ObjectMeta.Name, + Version: fullDS.GroupVersionKind().Version, + }, + } -// } -// } + return dsConnection, nil +} diff --git a/pkg/services/datasources/service/client/client_mock.go b/pkg/services/datasources/service/client/client_mock.go new file mode 100644 index 00000000000..03728d36e15 --- /dev/null +++ b/pkg/services/datasources/service/client/client_mock.go @@ -0,0 +1,98 @@ +// Code generated by mockery v2.53.3. DO NOT EDIT. + +package client + +import ( + context "context" + + v0alpha1 "github.com/grafana/grafana/pkg/apis/query/v0alpha1" + mock "github.com/stretchr/testify/mock" +) + +// MockDataSourceConnectionClient is an autogenerated mock type for the DataSourceConnectionClient type +type MockDataSourceConnectionClient struct { + mock.Mock +} + +type MockDataSourceConnectionClient_Expecter struct { + mock *mock.Mock +} + +func (_m *MockDataSourceConnectionClient) EXPECT() *MockDataSourceConnectionClient_Expecter { + return &MockDataSourceConnectionClient_Expecter{mock: &_m.Mock} +} + +// Get provides a mock function with given fields: ctx, group, version, name +func (_m *MockDataSourceConnectionClient) Get(ctx context.Context, group string, version string, name string) (*v0alpha1.DataSourceConnection, error) { + ret := _m.Called(ctx, group, version, name) + + if len(ret) == 0 { + panic("no return value specified for Get") + } + + var r0 *v0alpha1.DataSourceConnection + var r1 error + if rf, ok := ret.Get(0).(func(context.Context, string, string, string) (*v0alpha1.DataSourceConnection, error)); ok { + return rf(ctx, group, version, name) + } + if rf, ok := ret.Get(0).(func(context.Context, string, string, string) *v0alpha1.DataSourceConnection); ok { + r0 = rf(ctx, group, version, name) + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(*v0alpha1.DataSourceConnection) + } + } + + if rf, ok := ret.Get(1).(func(context.Context, string, string, string) error); ok { + r1 = rf(ctx, group, version, name) + } else { + r1 = ret.Error(1) + } + + return r0, r1 +} + +// MockDataSourceConnectionClient_Get_Call is a *mock.Call that shadows Run/Return methods with type explicit version for method 'Get' +type MockDataSourceConnectionClient_Get_Call struct { + *mock.Call +} + +// Get is a helper method to define mock.On call +// - ctx context.Context +// - group string +// - version string +// - name string +func (_e *MockDataSourceConnectionClient_Expecter) Get(ctx interface{}, group interface{}, version interface{}, name interface{}) *MockDataSourceConnectionClient_Get_Call { + return &MockDataSourceConnectionClient_Get_Call{Call: _e.mock.On("Get", ctx, group, version, name)} +} + +func (_c *MockDataSourceConnectionClient_Get_Call) Run(run func(ctx context.Context, group string, version string, name string)) *MockDataSourceConnectionClient_Get_Call { + _c.Call.Run(func(args mock.Arguments) { + run(args[0].(context.Context), args[1].(string), args[2].(string), args[3].(string)) + }) + return _c +} + +func (_c *MockDataSourceConnectionClient_Get_Call) Return(_a0 *v0alpha1.DataSourceConnection, _a1 error) *MockDataSourceConnectionClient_Get_Call { + _c.Call.Return(_a0, _a1) + return _c +} + +func (_c *MockDataSourceConnectionClient_Get_Call) RunAndReturn(run func(context.Context, string, string, string) (*v0alpha1.DataSourceConnection, error)) *MockDataSourceConnectionClient_Get_Call { + _c.Call.Return(run) + return _c +} + +// NewMockDataSourceConnectionClient creates a new instance of MockDataSourceConnectionClient. It also registers a testing interface on the mock and a cleanup function to assert the mocks expectations. +// The first argument is typically a *testing.T value. +func NewMockDataSourceConnectionClient(t interface { + mock.TestingT + Cleanup(func()) +}) *MockDataSourceConnectionClient { + mock := &MockDataSourceConnectionClient{} + mock.Mock.Test(t) + + t.Cleanup(func() { mock.AssertExpectations(t) }) + + return mock +}