diff --git a/.github/CODEOWNERS b/.github/CODEOWNERS index a6b9b13ae84..23d5b5efb2e 100644 --- a/.github/CODEOWNERS +++ b/.github/CODEOWNERS @@ -1137,6 +1137,8 @@ eslint-suppressions.json @grafanabot # Feature toggles /pkg/services/featuremgmt/ @grafana/grafana-backend-services-squad +# Data source migrations +/pkg/services/promtypemigration/ @grafana/partner-datasources @grafana/aws-datasources # Kind definitions /kinds/dashboard @grafana/dashboards-squad diff --git a/packages/grafana-data/src/types/featureToggles.gen.ts b/packages/grafana-data/src/types/featureToggles.gen.ts index 92e443e4867..614feba90a6 100644 --- a/packages/grafana-data/src/types/featureToggles.gen.ts +++ b/packages/grafana-data/src/types/featureToggles.gen.ts @@ -1132,4 +1132,9 @@ export interface FeatureToggles { * @default false */ azureResourcePickerUpdates?: boolean; + /** + * Checks for deprecated Prometheus authentication methods (SigV4 and Azure), installs the relevant data source, and migrates the Prometheus data sources + * @default false + */ + prometheusTypeMigration?: boolean; } diff --git a/pkg/registry/backgroundsvcs/background_services.go b/pkg/registry/backgroundsvcs/background_services.go index 9532b158b6c..2b66840cf05 100644 --- a/pkg/registry/backgroundsvcs/background_services.go +++ b/pkg/registry/backgroundsvcs/background_services.go @@ -33,6 +33,7 @@ import ( "github.com/grafana/grafana/pkg/services/pluginsintegration/pluginexternal" "github.com/grafana/grafana/pkg/services/pluginsintegration/plugininstaller" pluginStore "github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore" + "github.com/grafana/grafana/pkg/services/promtypemigration" "github.com/grafana/grafana/pkg/services/provisioning" publicdashboardsmetric "github.com/grafana/grafana/pkg/services/publicdashboards/metric" "github.com/grafana/grafana/pkg/services/rendering" @@ -71,6 +72,7 @@ func ProvideBackgroundServiceRegistry( pluginDashboardUpdater *plugindashboardsservice.DashboardUpdater, dashboardServiceImpl *service.DashboardServiceImpl, secretsGarbageCollectionWorker *secretsgarbagecollectionworker.Worker, + promTypeMigrationProvider promtypemigration.PromTypeMigrationProvider, // Need to make sure these are initialized, is there a better place to put them? _ dashboardsnapshots.Service, _ serviceaccounts.Service, @@ -118,6 +120,7 @@ func ProvideBackgroundServiceRegistry( pluginDashboardUpdater, dashboardServiceImpl, secretsGarbageCollectionWorker, + promTypeMigrationProvider, ) } diff --git a/pkg/server/wire.go b/pkg/server/wire.go index 5425639cf23..8dd9b0f9d9c 100644 --- a/pkg/server/wire.go +++ b/pkg/server/wire.go @@ -125,6 +125,7 @@ import ( pluginDashboards "github.com/grafana/grafana/pkg/services/pluginsintegration/dashboards" "github.com/grafana/grafana/pkg/services/pluginsintegration/pluginaccesscontrol" "github.com/grafana/grafana/pkg/services/preference/prefimpl" + promTypeMigration "github.com/grafana/grafana/pkg/services/promtypemigration" "github.com/grafana/grafana/pkg/services/publicdashboards" publicdashboardsApi "github.com/grafana/grafana/pkg/services/publicdashboards/api" publicdashboardsStore "github.com/grafana/grafana/pkg/services/publicdashboards/database" @@ -392,6 +393,10 @@ var wireBasicSet = wire.NewSet( secretsMigrations.ProvideDataSourceMigrationService, secretsMigrations.ProvideSecretMigrationProvider, wire.Bind(new(secretsMigrations.SecretMigrationProvider), new(*secretsMigrations.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)), diff --git a/pkg/server/wire_gen.go b/pkg/server/wire_gen.go index 6e0eb258d92..1002aadbfb8 100644 --- a/pkg/server/wire_gen.go +++ b/pkg/server/wire_gen.go @@ -194,6 +194,7 @@ import ( "github.com/grafana/grafana/pkg/services/pluginsintegration/sandbox" "github.com/grafana/grafana/pkg/services/pluginsintegration/serviceregistration" "github.com/grafana/grafana/pkg/services/preference/prefimpl" + "github.com/grafana/grafana/pkg/services/promtypemigration" "github.com/grafana/grafana/pkg/services/provisioning" "github.com/grafana/grafana/pkg/services/publicdashboards" api2 "github.com/grafana/grafana/pkg/services/publicdashboards/api" @@ -624,7 +625,10 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api if err != nil { return nil, err } - provisioningServiceImpl, err := provisioning.ProvideService(accessControl, cfg, sqlStore, pluginstoreService, dBstore, serviceService, notificationService, dashboardProvisioningService, service15, correlationsService, dashboardService, folderimplService, service13, searchService, quotaService, secretsService, orgService, receiverPermissionsService, tracingService, dualwriteService) + azurePromMigrationService := promtypemigration.ProvideAzurePromMigrationService(service15, inMemory, repoManager, pluginInstaller, cfg) + amazonPromMigrationService := promtypemigration.ProvideAmazonPromMigrationService(service15, inMemory, repoManager, pluginInstaller, cfg) + promTypeMigrationProviderImpl := promtypemigration.ProvidePromTypeMigrationProvider(serverLockService, featureToggles, azurePromMigrationService, amazonPromMigrationService) + provisioningServiceImpl, err := provisioning.ProvideService(accessControl, cfg, sqlStore, pluginstoreService, dBstore, serviceService, notificationService, dashboardProvisioningService, service15, correlationsService, dashboardService, folderimplService, service13, searchService, quotaService, secretsService, orgService, receiverPermissionsService, tracingService, dualwriteService, promTypeMigrationProviderImpl) if err != nil { return nil, err } @@ -857,7 +861,7 @@ func Initialize(ctx context.Context, cfg *setting.Cfg, opts Options, apiOpts api } ossUserProtectionImpl := authinfoimpl.ProvideOSSUserProtectionService() registration := authnimpl.ProvideRegistration(cfg, authnService, orgService, userAuthTokenService, acimplService, permissionRegistry, apikeyService, userService, authService, ossUserProtectionImpl, loginattemptimplService, quotaService, authinfoimplService, renderingService, featureToggles, oauthtokenService, socialService, remoteCache, ldapImpl, ossImpl, tracingService, tempuserService, notificationService) - backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, worker, serviceImpl, serviceAccountsProxy, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration) + backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, worker, promTypeMigrationProviderImpl, serviceImpl, serviceAccountsProxy, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration) usageStatsProvidersRegistry := usagestatssvcs.ProvideUsageStatsProvidersRegistry(acimplService, userService) server, err := New(opts, cfg, httpServer, acimplService, provisioningServiceImpl, backgroundServiceRegistry, usageStatsProvidersRegistry, statscollectorService, registerer) if err != nil { @@ -1202,7 +1206,10 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac if err != nil { return nil, err } - provisioningServiceImpl, err := provisioning.ProvideService(accessControl, cfg, sqlStore, pluginstoreService, dBstore, serviceService, notificationService, dashboardProvisioningService, service15, correlationsService, dashboardService, folderimplService, service13, searchService, quotaService, secretsService, orgService, receiverPermissionsService, tracingService, dualwriteService) + azurePromMigrationService := promtypemigration.ProvideAzurePromMigrationService(service15, inMemory, repoManager, pluginInstaller, cfg) + amazonPromMigrationService := promtypemigration.ProvideAmazonPromMigrationService(service15, inMemory, repoManager, pluginInstaller, cfg) + promTypeMigrationProviderImpl := promtypemigration.ProvidePromTypeMigrationProvider(serverLockService, featureToggles, azurePromMigrationService, amazonPromMigrationService) + provisioningServiceImpl, err := provisioning.ProvideService(accessControl, cfg, sqlStore, pluginstoreService, dBstore, serviceService, notificationService, dashboardProvisioningService, service15, correlationsService, dashboardService, folderimplService, service13, searchService, quotaService, secretsService, orgService, receiverPermissionsService, tracingService, dualwriteService, promTypeMigrationProviderImpl) if err != nil { return nil, err } @@ -1442,7 +1449,7 @@ func InitializeForTest(ctx context.Context, t sqlutil.ITestDB, testingT interfac } ossUserProtectionImpl := authinfoimpl.ProvideOSSUserProtectionService() registration := authnimpl.ProvideRegistration(cfg, authnService, orgService, userAuthTokenService, acimplService, permissionRegistry, apikeyService, userService, authService, ossUserProtectionImpl, loginattemptimplService, quotaService, authinfoimplService, renderingService, featureToggles, oauthtokentestService, socialService, remoteCache, ldapImpl, ossImpl, tracingService, tempuserService, notificationServiceMock) - backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, worker, serviceImpl, serviceAccountsProxy, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration) + backgroundServiceRegistry := backgroundsvcs.ProvideBackgroundServiceRegistry(httpServer, alertNG, cleanUpService, grafanaLive, gateway, notificationService, pluginstoreService, renderingService, userAuthTokenService, tracingService, provisioningServiceImpl, usageStats, statscollectorService, grafanaService, pluginsService, internalMetricsService, secretsService, remoteCache, storageService, searchService, entityEventsService, serviceAccountsService, grpcserverProvider, secretMigrationProviderImpl, loginattemptimplService, supportbundlesimplService, metricService, keyRetriever, angulardetectorsproviderDynamic, apiserverService, anonDeviceService, ssosettingsimplService, pluginexternalService, plugininstallerService, zanzanaReconciler, appregistryService, dashboardUpdater, dashboardServiceImpl, worker, promTypeMigrationProviderImpl, serviceImpl, serviceAccountsProxy, healthService, reflectionService, apiService, apiregistryService, idimplService, teamAPI, ssosettingsimplService, cloudmigrationService, registration) usageStatsProvidersRegistry := usagestatssvcs.ProvideUsageStatsProvidersRegistry(acimplService, userService) server, err := New(opts, cfg, httpServer, acimplService, provisioningServiceImpl, backgroundServiceRegistry, usageStatsProvidersRegistry, statscollectorService, registerer) if err != nil { @@ -1633,7 +1640,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, legacy.ProvideLegacyMigrator, 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, 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)), folderimpl.ProvideStore, wire.Bind(new(folder.Store), new(*folderimpl.FolderStoreImpl)), folderimpl.ProvideDashboardFolderStore, wire.Bind(new(folder.FolderStore), new(*folderimpl.DashboardFolderStoreImpl)), 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)), migrations2.ProvideDataSourceMigrationService, migrations2.ProvideSecretMigrationProvider, wire.Bind(new(migrations2.SecretMigrationProvider), new(*migrations2.SecretMigrationProviderImpl)), resourcepermissions.NewActionSetService, wire.Bind(new(accesscontrol.ActionResolver), new(resourcepermissions.ActionSetService)), wire.Bind(new(pluginaccesscontrol.ActionSetRegistry), new(resourcepermissions.ActionSetService)), permreg.ProvidePermissionRegistry, acimpl.ProvideAccessControl, 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, userimpl.ProvideVerifier, connectors.ProvideOrgRoleMapper, wire.Bind(new(user.Verifier), new(*userimpl.Verifier)), authz.WireSet, metadata.ProvideSecureValueMetadataStorage, metadata.ProvideKeeperMetadataStorage, metadata.ProvideDecryptStorage, decrypt.ProvideDecryptAuthorizer, decrypt.ProvideDecryptService, inline.ProvideInlineSecureValueService, encryption.ProvideDataKeyStorage, encryption.ProvideGlobalDataKeyStorage, encryption.ProvideEncryptedValueStorage, encryption.ProvideGlobalEncryptedValueStorage, service5.ProvideSecureValueService, validator.ProvideKeeperValidator, validator.ProvideSecureValueValidator, mutator.ProvideKeeperMutator, mutator.ProvideSecureValueMutator, migrator2.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, 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, legacy.ProvideLegacyMigrator, 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, 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)), folderimpl.ProvideStore, wire.Bind(new(folder.Store), new(*folderimpl.FolderStoreImpl)), folderimpl.ProvideDashboardFolderStore, wire.Bind(new(folder.FolderStore), new(*folderimpl.DashboardFolderStoreImpl)), 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)), migrations2.ProvideDataSourceMigrationService, migrations2.ProvideSecretMigrationProvider, wire.Bind(new(migrations2.SecretMigrationProvider), new(*migrations2.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, 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, userimpl.ProvideVerifier, connectors.ProvideOrgRoleMapper, wire.Bind(new(user.Verifier), new(*userimpl.Verifier)), authz.WireSet, metadata.ProvideSecureValueMetadataStorage, metadata.ProvideKeeperMetadataStorage, metadata.ProvideDecryptStorage, decrypt.ProvideDecryptAuthorizer, decrypt.ProvideDecryptService, inline.ProvideInlineSecureValueService, encryption.ProvideDataKeyStorage, encryption.ProvideGlobalDataKeyStorage, encryption.ProvideEncryptedValueStorage, encryption.ProvideGlobalEncryptedValueStorage, service5.ProvideSecureValueService, validator.ProvideKeeperValidator, validator.ProvideSecureValueValidator, mutator.ProvideKeeperMutator, mutator.ProvideSecureValueMutator, migrator2.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, apiserver.WireSet, apiregistry.WireSet, appregistry.WireSet, client.ProvideK8sClientWithFallback) 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/featuremgmt/registry.go b/pkg/services/featuremgmt/registry.go index 90cfeea5cab..a80b38250c6 100644 --- a/pkg/services/featuremgmt/registry.go +++ b/pkg/services/featuremgmt/registry.go @@ -1964,6 +1964,14 @@ var ( Owner: grafanaPartnerPluginsSquad, Expression: "false", }, + { + Name: "prometheusTypeMigration", + Description: "Checks for deprecated Prometheus authentication methods (SigV4 and Azure), installs the relevant data source, and migrates the Prometheus data sources", + Stage: FeatureStageExperimental, + RequiresRestart: true, + Owner: grafanaPartnerPluginsSquad, + Expression: "false", + }, } ) diff --git a/pkg/services/featuremgmt/toggles_gen.csv b/pkg/services/featuremgmt/toggles_gen.csv index faa23982f1b..2c1bae54ff0 100644 --- a/pkg/services/featuremgmt/toggles_gen.csv +++ b/pkg/services/featuremgmt/toggles_gen.csv @@ -252,3 +252,4 @@ teamFolders,experimental,@grafana/grafana-search-navigate-organise,false,false,f alertingTriage,experimental,@grafana/alerting-squad,false,false,true graphiteBackendMode,privatePreview,@grafana/partner-datasources,false,false,false azureResourcePickerUpdates,preview,@grafana/partner-datasources,false,false,true +prometheusTypeMigration,experimental,@grafana/partner-datasources,false,true,false diff --git a/pkg/services/featuremgmt/toggles_gen.go b/pkg/services/featuremgmt/toggles_gen.go index af44154ebeb..91761a72af7 100644 --- a/pkg/services/featuremgmt/toggles_gen.go +++ b/pkg/services/featuremgmt/toggles_gen.go @@ -1018,4 +1018,8 @@ const ( // FlagAzureResourcePickerUpdates // Enables the updated Azure Monitor resource picker FlagAzureResourcePickerUpdates = "azureResourcePickerUpdates" + + // FlagPrometheusTypeMigration + // Checks for deprecated Prometheus authentication methods (SigV4 and Azure), installs the relevant data source, and migrates the Prometheus data sources + FlagPrometheusTypeMigration = "prometheusTypeMigration" ) diff --git a/pkg/services/featuremgmt/toggles_gen.json b/pkg/services/featuremgmt/toggles_gen.json index 94ecf16c3f1..193ef690009 100644 --- a/pkg/services/featuremgmt/toggles_gen.json +++ b/pkg/services/featuremgmt/toggles_gen.json @@ -2714,6 +2714,23 @@ "frontend": true } }, + { + "metadata": { + "name": "prometheusTypeMigration", + "resourceVersion": "1757089774247", + "creationTimestamp": "2025-08-25T21:53:16Z", + "annotations": { + "grafana.app/updatedTimestamp": "2025-09-05 16:29:34.247055837 +0000 UTC" + } + }, + "spec": { + "description": "Checks for deprecated Prometheus authentication methods (SigV4 and Azure), installs the relevant data source, and migrates the Prometheus data sources", + "stage": "experimental", + "codeowner": "@grafana/partner-datasources", + "requiresRestart": true, + "expression": "false" + } + }, { "metadata": { "name": "provisioning", diff --git a/pkg/services/promtypemigration/amazon_prom_mig.go b/pkg/services/promtypemigration/amazon_prom_mig.go new file mode 100644 index 00000000000..33a9a35a35f --- /dev/null +++ b/pkg/services/promtypemigration/amazon_prom_mig.go @@ -0,0 +1,62 @@ +package promtypemigration + +import ( + "context" + + "github.com/grafana/grafana/pkg/plugins" + "github.com/grafana/grafana/pkg/plugins/manager/registry" + "github.com/grafana/grafana/pkg/plugins/repo" + "github.com/grafana/grafana/pkg/services/datasources" + "github.com/grafana/grafana/pkg/setting" +) + +type AmazonPromMigrationService struct { + promMigrationService +} + +func ProvideAmazonPromMigrationService( + dataSourcesService datasources.DataSourceService, + pluginRegistry registry.Service, + pluginRepo repo.Service, + pluginInstaller plugins.Installer, + cfg *setting.Cfg, +) *AmazonPromMigrationService { + return &AmazonPromMigrationService{ + promMigrationService: promMigrationService{ + dataSourcesService: dataSourcesService, + pluginRegistry: pluginRegistry, + pluginRepo: pluginRepo, + pluginInstaller: pluginInstaller, + cfg: cfg, + }, + } +} + +func (s *AmazonPromMigrationService) getPrometheusDataSources(ctx context.Context) ([]*datasources.DataSource, error) { + amazonPromDs := []*datasources.DataSource{} + query := &datasources.GetDataSourcesByTypeQuery{ + Type: datasources.DS_PROMETHEUS, + } + dsList, err := s.dataSourcesService.GetDataSourcesByType(ctx, query) + if err != nil { + return nil, err + } + for _, ds := range dsList { + if sigV4Auth, found := ds.JsonData.CheckGet("sigV4Auth"); found { + if enabled, err := sigV4Auth.Bool(); err != nil || !enabled { + continue + } + amazonPromDs = append(amazonPromDs, ds) + continue + } + } + return amazonPromDs, nil +} + +func (s *AmazonPromMigrationService) Migrate(ctx context.Context) error { + pds, err := s.getPrometheusDataSources(ctx) + if err != nil { + return err + } + return s.applyMigration(ctx, datasources.DS_AMAZON_PROMETHEUS, pds) +} diff --git a/pkg/services/promtypemigration/amazon_prom_mig_test.go b/pkg/services/promtypemigration/amazon_prom_mig_test.go new file mode 100644 index 00000000000..36ee7d5f06c --- /dev/null +++ b/pkg/services/promtypemigration/amazon_prom_mig_test.go @@ -0,0 +1,81 @@ +package promtypemigration + +import ( + "context" + "errors" + "testing" + + "github.com/grafana/grafana/pkg/components/simplejson" + "github.com/grafana/grafana/pkg/services/datasources" + "github.com/stretchr/testify/assert" +) + +func TestGetPrometheusDataSources_Amazon_ReturnsOnlyAmazonPrometheus(t *testing.T) { + ds1 := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{ + "sigV4Auth": true, + }), + } + ds2 := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{ + "sigV4Auth": false, + }), + } + ds3 := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{ + "sigV4Auth": true, + }), + } + ds4 := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{ + "sigV4Auth": nil, + }), + } + mock := &mockDataSourcesService{ + dataSources: []*datasources.DataSource{ds1, ds2, ds3, ds4}, + } + svc := &AmazonPromMigrationService{ + promMigrationService: promMigrationService{ + dataSourcesService: mock, + }, + } + + got, err := svc.getPrometheusDataSources(context.Background()) + assert.NoError(t, err) + assert.Len(t, got, 2) + assert.Contains(t, got, ds1) + assert.Contains(t, got, ds3) +} + +func TestGetPrometheusDataSources_Amazon_ErrorFromService(t *testing.T) { + mockSvc := &mockDataSourcesService{ + err: errors.New("service error"), + } + svc := &AmazonPromMigrationService{ + promMigrationService: promMigrationService{ + dataSourcesService: mockSvc, + }, + } + + got, err := svc.getPrometheusDataSources(context.Background()) + assert.Error(t, err) + assert.Nil(t, got) +} + +func TestGetPrometheusDataSources_Amazon_NoSigV4Auth(t *testing.T) { + ds := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{}), + } + mockSvc := &mockDataSourcesService{ + dataSources: []*datasources.DataSource{ds}, + } + svc := &AmazonPromMigrationService{ + promMigrationService: promMigrationService{ + dataSourcesService: mockSvc, + }, + } + + got, err := svc.getPrometheusDataSources(context.Background()) + assert.NoError(t, err) + assert.Empty(t, got) +} diff --git a/pkg/services/promtypemigration/azure_prom_mig.go b/pkg/services/promtypemigration/azure_prom_mig.go new file mode 100644 index 00000000000..0e8925a1352 --- /dev/null +++ b/pkg/services/promtypemigration/azure_prom_mig.go @@ -0,0 +1,63 @@ +package promtypemigration + +import ( + "context" + + "github.com/grafana/grafana/pkg/plugins" + "github.com/grafana/grafana/pkg/plugins/manager/registry" + "github.com/grafana/grafana/pkg/plugins/repo" + "github.com/grafana/grafana/pkg/services/datasources" + "github.com/grafana/grafana/pkg/setting" +) + +type AzurePromMigrationService struct { + promMigrationService +} + +func ProvideAzurePromMigrationService( + dataSourcesService datasources.DataSourceService, + pluginRegistry registry.Service, + pluginRepo repo.Service, + pluginInstaller plugins.Installer, + cfg *setting.Cfg, +) *AzurePromMigrationService { + return &AzurePromMigrationService{ + promMigrationService: promMigrationService{ + dataSourcesService: dataSourcesService, + pluginRegistry: pluginRegistry, + pluginRepo: pluginRepo, + pluginInstaller: pluginInstaller, + cfg: cfg, + }, + } +} + +func (s *AzurePromMigrationService) getPrometheusDataSources(ctx context.Context) ([]*datasources.DataSource, error) { + azurePromDs := []*datasources.DataSource{} + query := &datasources.GetDataSourcesByTypeQuery{ + Type: datasources.DS_PROMETHEUS, + } + dsList, err := s.dataSourcesService.GetDataSourcesByType(ctx, query) + if err != nil { + return nil, err + } + for _, ds := range dsList { + if azureAuth, found := ds.JsonData.CheckGet("azureCredentials"); found { + var val any + if val, err = azureAuth.Value(); err != nil || val == nil { + continue + } + azurePromDs = append(azurePromDs, ds) + continue + } + } + return azurePromDs, nil +} + +func (s *AzurePromMigrationService) Migrate(ctx context.Context) error { + pds, err := s.getPrometheusDataSources(ctx) + if err != nil { + return err + } + return s.applyMigration(ctx, datasources.DS_AZURE_PROMETHEUS, pds) +} diff --git a/pkg/services/promtypemigration/azure_prom_mig_test.go b/pkg/services/promtypemigration/azure_prom_mig_test.go new file mode 100644 index 00000000000..8bdb1402ff3 --- /dev/null +++ b/pkg/services/promtypemigration/azure_prom_mig_test.go @@ -0,0 +1,77 @@ +package promtypemigration + +import ( + "context" + "errors" + "testing" + + "github.com/grafana/grafana/pkg/components/simplejson" + "github.com/grafana/grafana/pkg/services/datasources" + "github.com/stretchr/testify/assert" +) + +func TestGetPrometheusDataSources_Azure_ReturnsOnlyAzurePrometheus(t *testing.T) { + ds1 := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{ + "azureCredentials": []any{}, + }), + } + ds2 := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{}), + } + ds3 := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{ + "azureCredentials": []any{}, + }), + } + ds4 := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{}), + } + mock := &mockDataSourcesService{ + dataSources: []*datasources.DataSource{ds1, ds2, ds3, ds4}, + } + svc := &AzurePromMigrationService{ + promMigrationService: promMigrationService{ + dataSourcesService: mock, + }, + } + + got, err := svc.getPrometheusDataSources(context.Background()) + assert.NoError(t, err) + assert.Len(t, got, 2) + assert.Contains(t, got, ds1) + assert.Contains(t, got, ds3) +} + +func TestGetPrometheusDataSources_Azure_ErrorFromService(t *testing.T) { + mockSvc := &mockDataSourcesService{ + err: errors.New("service error"), + } + svc := &AzurePromMigrationService{ + promMigrationService: promMigrationService{ + dataSourcesService: mockSvc, + }, + } + + got, err := svc.getPrometheusDataSources(context.Background()) + assert.Error(t, err) + assert.Nil(t, got) +} + +func TestGetPrometheusDataSources_Azure_NoAzureAuth(t *testing.T) { + ds := &datasources.DataSource{ + JsonData: simplejson.NewFromAny(map[string]any{}), + } + mockSvc := &mockDataSourcesService{ + dataSources: []*datasources.DataSource{ds}, + } + svc := &AzurePromMigrationService{ + promMigrationService: promMigrationService{ + dataSourcesService: mockSvc, + }, + } + + got, err := svc.getPrometheusDataSources(context.Background()) + assert.NoError(t, err) + assert.Empty(t, got) +} diff --git a/pkg/services/promtypemigration/migrator.go b/pkg/services/promtypemigration/migrator.go new file mode 100644 index 00000000000..abdc80fc7e7 --- /dev/null +++ b/pkg/services/promtypemigration/migrator.go @@ -0,0 +1,70 @@ +package promtypemigration + +import ( + "context" + "reflect" + "time" + + "github.com/grafana/grafana/pkg/infra/log" + "github.com/grafana/grafana/pkg/infra/serverlock" + "github.com/grafana/grafana/pkg/registry" + "github.com/grafana/grafana/pkg/services/featuremgmt" +) + +var logger = log.New("promds.migration") + +const actionName = "prom type migration task" + +type PromTypeMigrationService interface { + Migrate(ctx context.Context) error +} + +type PromTypeMigrationProvider interface { + registry.BackgroundService +} + +type PromTypeMigrationProviderImpl struct { + services []PromTypeMigrationService + features featuremgmt.FeatureToggles + ServerLockService *serverlock.ServerLockService +} + +func ProvidePromTypeMigrationProvider( + serverLockService *serverlock.ServerLockService, + features featuremgmt.FeatureToggles, + promAzureAuthMigrationService *AzurePromMigrationService, + promAmazonAuthMigrationService *AmazonPromMigrationService, +) *PromTypeMigrationProviderImpl { + return &PromTypeMigrationProviderImpl{ + ServerLockService: serverLockService, + features: features, + services: []PromTypeMigrationService{promAzureAuthMigrationService, promAmazonAuthMigrationService}, + } +} + +func (s *PromTypeMigrationProviderImpl) Run(ctx context.Context) error { + if !s.features.IsEnabled(ctx, featuremgmt.FlagPrometheusTypeMigration) { + return nil + } + return s.migrate(ctx) +} + +// migrate Run migration services. This will block until all services have exited. +// This should only be called once at startup +func (s *PromTypeMigrationProviderImpl) migrate(ctx context.Context) error { + err := s.ServerLockService.LockExecuteAndRelease(ctx, actionName, time.Minute*10, func(context.Context) { + for _, service := range s.services { + serviceName := reflect.TypeOf(service).String() + logger.Debug("Starting prom data source type migration service", "service", serviceName) + err := service.Migrate(ctx) + if err != nil { + logger.Error("Stopped prom data source type migration service", "service", serviceName, "reason", err) + } + logger.Debug("Finished prom data source type migration service", "service", serviceName) + } + }) + if err != nil { + logger.Error("Server lock for prom data source type migration already exists") + } + return nil +} diff --git a/pkg/services/promtypemigration/prom_mig.go b/pkg/services/promtypemigration/prom_mig.go new file mode 100644 index 00000000000..63b08f17d8d --- /dev/null +++ b/pkg/services/promtypemigration/prom_mig.go @@ -0,0 +1,83 @@ +package promtypemigration + +import ( + "context" + "runtime" + + "github.com/grafana/grafana/pkg/components/simplejson" + "github.com/grafana/grafana/pkg/plugins" + "github.com/grafana/grafana/pkg/plugins/manager/registry" + "github.com/grafana/grafana/pkg/plugins/repo" + "github.com/grafana/grafana/pkg/services/datasources" + "github.com/grafana/grafana/pkg/setting" +) + +type PromMigrationHandler interface { + Migrate(context.Context, *promMigrationService) error +} + +type promMigrationService struct { + cfg *setting.Cfg + dataSourcesService datasources.DataSourceService + pluginRegistry registry.Service + pluginRepo repo.Service + pluginInstaller plugins.Installer +} + +func (s *promMigrationService) applyMigration(ctx context.Context, pluginID string, promDataSources []*datasources.DataSource) error { + if len(promDataSources) == 0 { + return nil + } + + // check to see if prom is installed, if not install it + if _, installed := s.pluginRegistry.Plugin(ctx, pluginID, ""); !installed { + compatOpts := plugins.NewAddOpts(s.cfg.BuildVersion, runtime.GOOS, runtime.GOARCH, "") + err := s.pluginInstaller.Add(ctx, pluginID, "", compatOpts) + if err != nil { + return err + } + } + + logger.Debug("performing prometheus data source type migration", "plugin", pluginID) + + for _, ds := range promDataSources { + err := s.updateDataSourceType(ctx, ds, pluginID) + if err != nil { + return err + } + } + + logger.Debug("prometheus data source type migration complete", "plugin", pluginID) + + return nil +} + +func (s *promMigrationService) updateDataSourceType(ctx context.Context, ds *datasources.DataSource, newType string) error { + secureJsonData, err := s.dataSourcesService.DecryptedValues(ctx, ds) + if err != nil { + return err + } + if ds.JsonData == nil { + logger.Debug("no JsonData found", "data source ID", ds.ID) + ds.JsonData = &simplejson.Json{} + } + ds.JsonData.Set("prometheus-type-migration", true) + _, err = s.dataSourcesService.UpdateDataSource(ctx, &datasources.UpdateDataSourceCommand{ + ID: ds.ID, + Type: newType, + OrgID: ds.OrgID, + UID: ds.UID, + Name: ds.Name, + JsonData: ds.JsonData, + SecureJsonData: secureJsonData, + + // These are needed by the SQL function due to UseBool and MustCols + IsDefault: ds.IsDefault, + BasicAuth: ds.BasicAuth, + WithCredentials: ds.WithCredentials, + ReadOnly: ds.ReadOnly, + User: ds.User, + Database: ds.Database, + }) + return err +} diff --git a/pkg/services/promtypemigration/prom_mig_test.go b/pkg/services/promtypemigration/prom_mig_test.go new file mode 100644 index 00000000000..57bdac3423f --- /dev/null +++ b/pkg/services/promtypemigration/prom_mig_test.go @@ -0,0 +1,120 @@ +package promtypemigration + +import ( + "context" + "errors" + "testing" + + "github.com/grafana/grafana/pkg/components/simplejson" + "github.com/grafana/grafana/pkg/plugins" + "github.com/grafana/grafana/pkg/services/datasources" + "github.com/grafana/grafana/pkg/setting" + "github.com/stretchr/testify/assert" +) + +// Mocks + +type mockPluginRegistry struct { + installed bool +} + +func (m *mockPluginRegistry) Plugin(ctx context.Context, id string, _ string) (*plugins.Plugin, bool) { + if m.installed { + return &plugins.Plugin{}, true + } + return &plugins.Plugin{}, false +} +func (m *mockPluginRegistry) Plugins(ctx context.Context) []*plugins.Plugin { return nil } +func (m *mockPluginRegistry) Add(ctx context.Context, plugin *plugins.Plugin) error { return nil } +func (m *mockPluginRegistry) Remove(ctx context.Context, id, version string) error { return nil } + +type mockPluginInstaller struct { + addCalled bool + addErr error +} + +func (m *mockPluginInstaller) Add(ctx context.Context, pluginID, version string, opts plugins.AddOpts) error { + m.addCalled = true + return m.addErr +} +func (m *mockPluginInstaller) Remove(ctx context.Context, pluginID, version string) error { + return nil +} + +type mockDataSourcesService struct { + datasources.DataSourceService + dataSources []*datasources.DataSource + err error +} + +func (m *mockDataSourcesService) DecryptedValues(ctx context.Context, ds *datasources.DataSource) (map[string]string, error) { + return map[string]string{}, nil +} + +func (m *mockDataSourcesService) UpdateDataSource(ctx context.Context, cmd *datasources.UpdateDataSourceCommand) (*datasources.DataSource, error) { + return &datasources.DataSource{}, m.err +} + +func (m *mockDataSourcesService) GetDataSourcesByType(ctx context.Context, query *datasources.GetDataSourcesByTypeQuery) ([]*datasources.DataSource, error) { + return m.dataSources, m.err +} + +// Test cases + +func TestApplyMigration_NoDataSources(t *testing.T) { + svc := &promMigrationService{} + err := svc.applyMigration(context.Background(), "prometheus", []*datasources.DataSource{}) + assert.NoError(t, err) +} + +func TestApplyMigration_PluginAlreadyInstalled(t *testing.T) { + ds := &datasources.DataSource{ID: 1, JsonData: simplejson.New()} + svc := &promMigrationService{ + cfg: &setting.Cfg{BuildVersion: "1.0"}, + dataSourcesService: &mockDataSourcesService{}, + pluginRegistry: &mockPluginRegistry{installed: true}, + pluginInstaller: &mockPluginInstaller{}, + } + err := svc.applyMigration(context.Background(), "prometheus", []*datasources.DataSource{ds}) + assert.NoError(t, err) +} + +func TestApplyMigration_PluginNotInstalled_InstallSucceeds(t *testing.T) { + ds := &datasources.DataSource{ID: 1, JsonData: simplejson.New()} + installer := &mockPluginInstaller{} + svc := &promMigrationService{ + cfg: &setting.Cfg{BuildVersion: "1.0"}, + dataSourcesService: &mockDataSourcesService{}, + pluginRegistry: &mockPluginRegistry{installed: false}, + pluginInstaller: installer, + } + err := svc.applyMigration(context.Background(), "prometheus", []*datasources.DataSource{ds}) + assert.NoError(t, err) + assert.True(t, installer.addCalled) +} + +func TestApplyMigration_PluginNotInstalled_InstallFails(t *testing.T) { + ds := &datasources.DataSource{ID: 1} + installer := &mockPluginInstaller{addErr: errors.New("install failed")} + svc := &promMigrationService{ + cfg: &setting.Cfg{BuildVersion: "1.0"}, + dataSourcesService: &mockDataSourcesService{}, + pluginRegistry: &mockPluginRegistry{installed: false}, + pluginInstaller: installer, + } + err := svc.applyMigration(context.Background(), "prometheus", []*datasources.DataSource{ds}) + assert.EqualError(t, err, "install failed") +} + +func TestApplyMigration_UpdateDataSourceFails(t *testing.T) { + ds := &datasources.DataSource{ID: 1, JsonData: simplejson.New()} + dataSvc := &mockDataSourcesService{err: errors.New("update failed")} + svc := &promMigrationService{ + cfg: &setting.Cfg{BuildVersion: "1.0"}, + dataSourcesService: dataSvc, + pluginRegistry: &mockPluginRegistry{installed: true}, + pluginInstaller: &mockPluginInstaller{}, + } + err := svc.applyMigration(context.Background(), "prometheus", []*datasources.DataSource{ds}) + assert.EqualError(t, err, "update failed") +} diff --git a/pkg/services/provisioning/provisioning.go b/pkg/services/provisioning/provisioning.go index 6972eb22213..3f831f9d4a3 100644 --- a/pkg/services/provisioning/provisioning.go +++ b/pkg/services/provisioning/provisioning.go @@ -27,6 +27,7 @@ import ( "github.com/grafana/grafana/pkg/services/org" "github.com/grafana/grafana/pkg/services/pluginsintegration/pluginsettings" "github.com/grafana/grafana/pkg/services/pluginsintegration/pluginstore" + "github.com/grafana/grafana/pkg/services/promtypemigration" prov_alerting "github.com/grafana/grafana/pkg/services/provisioning/alerting" "github.com/grafana/grafana/pkg/services/provisioning/dashboards" "github.com/grafana/grafana/pkg/services/provisioning/datasources" @@ -59,6 +60,7 @@ func ProvideService( resourcePermissions accesscontrol.ReceiverPermissionsService, tracer tracing.Tracer, dual dualwrite.Service, + promTypeMigrationProvider promtypemigration.PromTypeMigrationProvider, ) (*ProvisioningServiceImpl, error) { s := &ProvisioningServiceImpl{ Cfg: cfg, @@ -85,6 +87,7 @@ func ProvideService( folderService: folderService, resourcePermissions: resourcePermissions, tracer: tracer, + migratePrometheusType: promTypeMigrationProvider.Run, } if err := s.setDashboardProvisioner(); err != nil { @@ -120,6 +123,7 @@ func newProvisioningServiceImpl( newDashboardProvisioner dashboards.DashboardProvisionerFactory, provisionDatasources func(context.Context, string, datasources.BaseDataSourceService, datasources.CorrelationsStore, org.Service) error, provisionPlugins func(context.Context, string, pluginstore.Store, pluginsettings.Service, org.Service) error, + migratePrometheusType func(context.Context) error, searchService searchV2.SearchService, ) (*ProvisioningServiceImpl, error) { s := &ProvisioningServiceImpl{ @@ -129,6 +133,7 @@ func newProvisioningServiceImpl( provisionPlugins: provisionPlugins, Cfg: setting.NewCfg(), searchService: searchService, + migratePrometheusType: migratePrometheusType, } if err := s.setDashboardProvisioner(); err != nil { @@ -168,6 +173,7 @@ type ProvisioningServiceImpl struct { tracer tracing.Tracer dual dualwrite.Service onceInitProvisioners sync.Once + migratePrometheusType func(context.Context) error } func (ps *ProvisioningServiceImpl) RunInitProvisioners(ctx context.Context) error { @@ -201,6 +207,15 @@ func (ps *ProvisioningServiceImpl) Run(ctx context.Context) error { ps.log.Error("Failed to provision alerting", "error", err) return } + + // Migrating prom types relies on data source provisioning to already be completed + // If we can make services depend on other services completing first, + // then we should remove this from provisioning + err = ps.migratePrometheusType(ctx) + if err != nil { + ps.log.Error("Failed to migrate Prometheus type", "error", err) + return + } }) if err != nil { diff --git a/pkg/services/provisioning/provisioning_test.go b/pkg/services/provisioning/provisioning_test.go index 6f9558d4702..ee4abc5029b 100644 --- a/pkg/services/provisioning/provisioning_test.go +++ b/pkg/services/provisioning/provisioning_test.go @@ -170,6 +170,9 @@ func setup(t *testing.T) *serviceTestStruct { func(context.Context, string, pluginstore.Store, pluginsettings.Service, org.Service) error { return nil }, + func(context.Context) error { + return nil + }, searchStub, ) service.provisionAlerting = func(context.Context, prov_alerting.ProvisionerConfig) error {