CloudMigrations: Move business logic out of api layer (#86406)
* move run migration to the cloudmigrationimpl layer * add migration run list logic down a layer * remove useless comments * pull cms calls into their own service
This commit is contained in:
@@ -1,10 +1,7 @@
|
|||||||
package api
|
package api
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
|
||||||
"encoding/json"
|
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
|
||||||
"net/http"
|
"net/http"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
|
||||||
@@ -172,7 +169,6 @@ func (cma *CloudMigrationAPI) CreateMigration(c *contextmodel.ReqContext) respon
|
|||||||
func (cma *CloudMigrationAPI) RunMigration(c *contextmodel.ReqContext) response.Response {
|
func (cma *CloudMigrationAPI) RunMigration(c *contextmodel.ReqContext) response.Response {
|
||||||
ctx, span := cma.tracer.Start(c.Req.Context(), "MigrationAPI.RunMigration")
|
ctx, span := cma.tracer.Start(c.Req.Context(), "MigrationAPI.RunMigration")
|
||||||
defer span.End()
|
defer span.End()
|
||||||
logger := cma.log.FromContext(ctx)
|
|
||||||
|
|
||||||
stringID := web.Params(c.Req)[":id"]
|
stringID := web.Params(c.Req)[":id"]
|
||||||
id, err := strconv.ParseInt(stringID, 10, 64)
|
id, err := strconv.ParseInt(stringID, 10, 64)
|
||||||
@@ -180,70 +176,10 @@ func (cma *CloudMigrationAPI) RunMigration(c *contextmodel.ReqContext) response.
|
|||||||
return response.Error(http.StatusBadRequest, "id is invalid", err)
|
return response.Error(http.StatusBadRequest, "id is invalid", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Get migration to read the auth token
|
result, err := cma.cloudMigrationService.RunMigration(ctx, id)
|
||||||
migration, err := cma.cloudMigrationService.GetMigration(ctx, id)
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return response.Error(http.StatusInternalServerError, "migration get error", err)
|
return response.Error(http.StatusInternalServerError, "migration run error", err)
|
||||||
}
|
}
|
||||||
// get CMS path from the config
|
|
||||||
domain, err := cma.cloudMigrationService.ParseCloudMigrationConfig()
|
|
||||||
if err != nil {
|
|
||||||
return response.Error(http.StatusInternalServerError, "config parse error", err)
|
|
||||||
}
|
|
||||||
path := fmt.Sprintf("https://cms-%s.%s/cloud-migrations/api/v1/migrate-data", migration.ClusterSlug, domain)
|
|
||||||
|
|
||||||
// Get migration data JSON
|
|
||||||
body, err := cma.cloudMigrationService.GetMigrationDataJSON(ctx, id)
|
|
||||||
if err != nil {
|
|
||||||
cma.log.Error("error getting the json request body for migration run", "err", err.Error())
|
|
||||||
return response.Error(http.StatusInternalServerError, "migration data get error", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
req, err := http.NewRequest(http.MethodPost, path, bytes.NewReader(body))
|
|
||||||
if err != nil {
|
|
||||||
cma.log.Error("error creating http request for cloud migration run", "err", err.Error())
|
|
||||||
return response.Error(http.StatusInternalServerError, "http request error", err)
|
|
||||||
}
|
|
||||||
req.Header.Set("Content-Type", "application/json")
|
|
||||||
req.Header.Set("Authorization", fmt.Sprintf("Bearer %d:%s", migration.StackID, migration.AuthToken))
|
|
||||||
|
|
||||||
client := &http.Client{}
|
|
||||||
resp, err := client.Do(req)
|
|
||||||
if err != nil {
|
|
||||||
cma.log.Error("error sending http request for cloud migration run", "err", err.Error())
|
|
||||||
return response.Error(http.StatusInternalServerError, "http request error", err)
|
|
||||||
} else if resp.StatusCode >= 400 {
|
|
||||||
cma.log.Error("received error response for cloud migration run", "statusCode", resp.StatusCode)
|
|
||||||
return response.Error(http.StatusInternalServerError, "http request error", fmt.Errorf("http request error while migrating data"))
|
|
||||||
}
|
|
||||||
|
|
||||||
defer func() {
|
|
||||||
if err := resp.Body.Close(); err != nil {
|
|
||||||
logger.Error("closing request body: %w", err)
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
// read response so we can unmarshal it
|
|
||||||
respData, err := io.ReadAll(resp.Body)
|
|
||||||
if err != nil {
|
|
||||||
logger.Error("reading response body: %w", err)
|
|
||||||
return response.Error(http.StatusInternalServerError, "reading migration run response", err)
|
|
||||||
}
|
|
||||||
var result cloudmigration.MigrateDataResponseDTO
|
|
||||||
if err := json.Unmarshal(respData, &result); err != nil {
|
|
||||||
logger.Error("unmarshalling response body: %w", err)
|
|
||||||
return response.Error(http.StatusInternalServerError, "unmarshalling migration run response", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
runID, err := cma.cloudMigrationService.SaveMigrationRun(ctx, &cloudmigration.CloudMigrationRun{
|
|
||||||
CloudMigrationUID: stringID,
|
|
||||||
Result: respData,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
response.Error(http.StatusInternalServerError, "migration run save error", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
result.RunID = runID
|
|
||||||
|
|
||||||
return response.JSON(http.StatusOK, result)
|
return response.JSON(http.StatusOK, result)
|
||||||
}
|
}
|
||||||
@@ -309,22 +245,9 @@ func (cma *CloudMigrationAPI) GetMigrationRunList(c *contextmodel.ReqContext) re
|
|||||||
ctx, span := cma.tracer.Start(c.Req.Context(), "MigrationAPI.GetMigrationRunList")
|
ctx, span := cma.tracer.Start(c.Req.Context(), "MigrationAPI.GetMigrationRunList")
|
||||||
defer span.End()
|
defer span.End()
|
||||||
|
|
||||||
migrationStatuses, err := cma.cloudMigrationService.GetMigrationStatusList(ctx, web.Params(c.Req)[":id"])
|
runList, err := cma.cloudMigrationService.GetMigrationRunList(ctx, web.Params(c.Req)[":id"])
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return response.Error(http.StatusInternalServerError, "migration status error", err)
|
return response.Error(http.StatusInternalServerError, "list migration status error", err)
|
||||||
}
|
|
||||||
|
|
||||||
runList := cloudmigration.CloudMigrationRunList{Runs: []cloudmigration.MigrateDataResponseDTO{}}
|
|
||||||
for _, s := range migrationStatuses {
|
|
||||||
// attempt to bind the raw result to a list of response item DTOs
|
|
||||||
r := cloudmigration.MigrateDataResponseDTO{
|
|
||||||
Items: []cloudmigration.MigrateDataResponseItemDTO{},
|
|
||||||
}
|
|
||||||
if err := json.Unmarshal(s.Result, &r); err != nil {
|
|
||||||
return response.Error(http.StatusInternalServerError, "error unmarshalling migration response items", err)
|
|
||||||
}
|
|
||||||
r.RunID = s.ID
|
|
||||||
runList.Runs = append(runList.Runs, r)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return response.JSON(http.StatusOK, runList)
|
return response.JSON(http.StatusOK, runList)
|
||||||
|
|||||||
@@ -7,16 +7,15 @@ import (
|
|||||||
type Service interface {
|
type Service interface {
|
||||||
CreateToken(context.Context) (CreateAccessTokenResponse, error)
|
CreateToken(context.Context) (CreateAccessTokenResponse, error)
|
||||||
ValidateToken(context.Context, CloudMigration) error
|
ValidateToken(context.Context, CloudMigration) error
|
||||||
// migration
|
|
||||||
GetMigration(context.Context, int64) (*CloudMigration, error)
|
|
||||||
GetMigrationList(context.Context) (*CloudMigrationListResponse, error)
|
|
||||||
CreateMigration(context.Context, CloudMigrationRequest) (*CloudMigrationResponse, error)
|
|
||||||
GetMigrationDataJSON(context.Context, int64) ([]byte, error)
|
|
||||||
UpdateMigration(context.Context, int64, CloudMigrationRequest) (*CloudMigrationResponse, error)
|
|
||||||
GetMigrationStatus(context.Context, string, string) (*CloudMigrationRun, error)
|
|
||||||
GetMigrationStatusList(context.Context, string) ([]*CloudMigrationRun, error)
|
|
||||||
DeleteMigration(context.Context, int64) (*CloudMigration, error)
|
|
||||||
SaveMigrationRun(context.Context, *CloudMigrationRun) (int64, error)
|
|
||||||
|
|
||||||
ParseCloudMigrationConfig() (string, error)
|
CreateMigration(context.Context, CloudMigrationRequest) (*CloudMigrationResponse, error)
|
||||||
|
GetMigration(context.Context, int64) (*CloudMigration, error)
|
||||||
|
DeleteMigration(context.Context, int64) (*CloudMigration, error)
|
||||||
|
UpdateMigration(context.Context, int64, CloudMigrationRequest) (*CloudMigrationResponse, error)
|
||||||
|
GetMigrationList(context.Context) (*CloudMigrationListResponse, error)
|
||||||
|
|
||||||
|
RunMigration(context.Context, int64) (*MigrateDataResponseDTO, error)
|
||||||
|
SaveMigrationRun(context.Context, *CloudMigrationRun) (int64, error)
|
||||||
|
GetMigrationStatus(context.Context, string, string) (*CloudMigrationRun, error)
|
||||||
|
GetMigrationRunList(context.Context, string) (*CloudMigrationRunList, error)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,20 +1,22 @@
|
|||||||
package cloudmigrationimpl
|
package cloudmigrationimpl
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
|
||||||
"context"
|
"context"
|
||||||
"encoding/base64"
|
"encoding/base64"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"strconv"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/grafana/grafana/pkg/api/response"
|
||||||
"github.com/grafana/grafana/pkg/api/routing"
|
"github.com/grafana/grafana/pkg/api/routing"
|
||||||
"github.com/grafana/grafana/pkg/infra/db"
|
"github.com/grafana/grafana/pkg/infra/db"
|
||||||
"github.com/grafana/grafana/pkg/infra/log"
|
"github.com/grafana/grafana/pkg/infra/log"
|
||||||
"github.com/grafana/grafana/pkg/infra/tracing"
|
"github.com/grafana/grafana/pkg/infra/tracing"
|
||||||
"github.com/grafana/grafana/pkg/services/cloudmigration"
|
"github.com/grafana/grafana/pkg/services/cloudmigration"
|
||||||
"github.com/grafana/grafana/pkg/services/cloudmigration/api"
|
"github.com/grafana/grafana/pkg/services/cloudmigration/api"
|
||||||
|
"github.com/grafana/grafana/pkg/services/cloudmigration/cmsclient"
|
||||||
"github.com/grafana/grafana/pkg/services/contexthandler"
|
"github.com/grafana/grafana/pkg/services/contexthandler"
|
||||||
"github.com/grafana/grafana/pkg/services/dashboards"
|
"github.com/grafana/grafana/pkg/services/dashboards"
|
||||||
"github.com/grafana/grafana/pkg/services/datasources"
|
"github.com/grafana/grafana/pkg/services/datasources"
|
||||||
@@ -33,7 +35,8 @@ type Service struct {
|
|||||||
log *log.ConcreteLogger
|
log *log.ConcreteLogger
|
||||||
cfg *setting.Cfg
|
cfg *setting.Cfg
|
||||||
|
|
||||||
features featuremgmt.FeatureToggles
|
features featuremgmt.FeatureToggles
|
||||||
|
cmsClient cmsclient.Client
|
||||||
|
|
||||||
dsService datasources.DataSourceService
|
dsService datasources.DataSourceService
|
||||||
gcomService gcom.Service
|
gcomService gcom.Service
|
||||||
@@ -70,9 +73,9 @@ func ProvideService(
|
|||||||
tracer tracing.Tracer,
|
tracer tracing.Tracer,
|
||||||
dashboardService dashboards.DashboardService,
|
dashboardService dashboards.DashboardService,
|
||||||
folderService folder.Service,
|
folderService folder.Service,
|
||||||
) cloudmigration.Service {
|
) (cloudmigration.Service, error) {
|
||||||
if !features.IsEnabledGlobally(featuremgmt.FlagOnPremToCloudMigrations) {
|
if !features.IsEnabledGlobally(featuremgmt.FlagOnPremToCloudMigrations) {
|
||||||
return &NoopServiceImpl{}
|
return &NoopServiceImpl{}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
s := &Service{
|
s := &Service{
|
||||||
@@ -90,11 +93,18 @@ func ProvideService(
|
|||||||
}
|
}
|
||||||
s.api = api.RegisterApi(routeRegister, s, tracer)
|
s.api = api.RegisterApi(routeRegister, s, tracer)
|
||||||
|
|
||||||
|
// get CMS path from the config
|
||||||
|
domain, err := s.parseCloudMigrationConfig()
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("config parse error: %w", err)
|
||||||
|
}
|
||||||
|
s.cmsClient = cmsclient.NewCMSClient(domain)
|
||||||
|
|
||||||
if err := s.registerMetrics(prom, s.metrics); err != nil {
|
if err := s.registerMetrics(prom, s.metrics); err != nil {
|
||||||
s.log.Warn("error registering prom metrics", "error", err.Error())
|
s.log.Warn("error registering prom metrics", "error", err.Error())
|
||||||
}
|
}
|
||||||
|
|
||||||
return s
|
return s, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) CreateToken(ctx context.Context) (cloudmigration.CreateAccessTokenResponse, error) {
|
func (s *Service) CreateToken(ctx context.Context) (cloudmigration.CreateAccessTokenResponse, error) {
|
||||||
@@ -213,44 +223,9 @@ func (s *Service) findAccessPolicyByName(ctx context.Context, regionSlug, access
|
|||||||
func (s *Service) ValidateToken(ctx context.Context, cm cloudmigration.CloudMigration) error {
|
func (s *Service) ValidateToken(ctx context.Context, cm cloudmigration.CloudMigration) error {
|
||||||
ctx, span := s.tracer.Start(ctx, "CloudMigrationService.ValidateToken")
|
ctx, span := s.tracer.Start(ctx, "CloudMigrationService.ValidateToken")
|
||||||
defer span.End()
|
defer span.End()
|
||||||
logger := s.log.FromContext(ctx)
|
|
||||||
|
|
||||||
// get CMS path from the config
|
if err := s.cmsClient.ValidateKey(ctx, cm); err != nil {
|
||||||
domain, err := s.ParseCloudMigrationConfig()
|
return fmt.Errorf("validating key: %w", err)
|
||||||
if err != nil {
|
|
||||||
return fmt.Errorf("config parse error: %w", err)
|
|
||||||
}
|
|
||||||
path := fmt.Sprintf("https://cms-%s.%s/cloud-migrations/api/v1/validate-key", cm.ClusterSlug, domain)
|
|
||||||
|
|
||||||
// validation is an empty POST to CMS with the authorization header included
|
|
||||||
req, err := http.NewRequest("POST", path, bytes.NewReader(nil))
|
|
||||||
if err != nil {
|
|
||||||
logger.Error("error creating http request for token validation", "err", err.Error())
|
|
||||||
return fmt.Errorf("http request error: %w", err)
|
|
||||||
}
|
|
||||||
req.Header.Set("Content-Type", "application/json")
|
|
||||||
req.Header.Set("Authorization", fmt.Sprintf("Bearer %d:%s", cm.StackID, cm.AuthToken))
|
|
||||||
|
|
||||||
client := &http.Client{}
|
|
||||||
resp, err := client.Do(req)
|
|
||||||
if err != nil {
|
|
||||||
logger.Error("error sending http request for token validation", "err", err.Error())
|
|
||||||
return fmt.Errorf("http request error: %w", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
defer func() {
|
|
||||||
if err := resp.Body.Close(); err != nil {
|
|
||||||
logger.Error("closing request body", "err", err.Error())
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
|
|
||||||
if resp.StatusCode != 200 {
|
|
||||||
var errResp map[string]any
|
|
||||||
if err := json.NewDecoder(resp.Body).Decode(&errResp); err != nil {
|
|
||||||
logger.Error("decoding error response", "err", err.Error())
|
|
||||||
} else {
|
|
||||||
return fmt.Errorf("token validation failure: %v", errResp)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
@@ -323,10 +298,52 @@ func (s *Service) UpdateMigration(ctx context.Context, id int64, cm cloudmigrati
|
|||||||
return nil, nil
|
return nil, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) GetMigrationDataJSON(ctx context.Context, id int64) ([]byte, error) {
|
func (s *Service) RunMigration(ctx context.Context, id int64) (*cloudmigration.MigrateDataResponseDTO, error) {
|
||||||
|
// Get migration to read the auth token
|
||||||
|
migration, err := s.GetMigration(ctx, id)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("migration get error: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Get migration data JSON
|
||||||
|
request, err := s.getMigrationDataJSON(ctx)
|
||||||
|
if err != nil {
|
||||||
|
s.log.Error("error getting the json request body for migration run", "err", err.Error())
|
||||||
|
return nil, fmt.Errorf("migration data get error: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Call the cms service
|
||||||
|
resp, err := s.cmsClient.MigrateData(ctx, *migration, *request)
|
||||||
|
if err != nil {
|
||||||
|
s.log.Error("error migrating data: %w", err)
|
||||||
|
return nil, fmt.Errorf("migrate data error: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TODO update cloud migration run schema to treat the result as a first-class citizen
|
||||||
|
respData, err := json.Marshal(resp)
|
||||||
|
if err != nil {
|
||||||
|
s.log.Error("error marshalling migration response data: %w", err)
|
||||||
|
return nil, fmt.Errorf("marshalling migration response data: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// save the result of the migration
|
||||||
|
runID, err := s.SaveMigrationRun(ctx, &cloudmigration.CloudMigrationRun{
|
||||||
|
CloudMigrationUID: strconv.Itoa(int(id)),
|
||||||
|
Result: respData,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
response.Error(http.StatusInternalServerError, "migration run save error", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
resp.RunID = runID
|
||||||
|
|
||||||
|
return resp, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Service) getMigrationDataJSON(ctx context.Context) (*cloudmigration.MigrateDataRequestDTO, error) {
|
||||||
var migrationDataSlice []cloudmigration.MigrateDataRequestItemDTO
|
var migrationDataSlice []cloudmigration.MigrateDataRequestItemDTO
|
||||||
// Data sources
|
// Data sources
|
||||||
dataSources, err := s.getDataSources(ctx, id)
|
dataSources, err := s.getDataSources(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.log.Error("Failed to get datasources", "err", err)
|
s.log.Error("Failed to get datasources", "err", err)
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -341,7 +358,7 @@ func (s *Service) GetMigrationDataJSON(ctx context.Context, id int64) ([]byte, e
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Dashboards
|
// Dashboards
|
||||||
dashboards, err := s.getDashboards(ctx, id)
|
dashboards, err := s.getDashboards(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.log.Error("Failed to get dashboards", "err", err)
|
s.log.Error("Failed to get dashboards", "err", err)
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -358,7 +375,7 @@ func (s *Service) GetMigrationDataJSON(ctx context.Context, id int64) ([]byte, e
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Folders
|
// Folders
|
||||||
folders, err := s.getFolders(ctx, id)
|
folders, err := s.getFolders(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.log.Error("Failed to get folders", "err", err)
|
s.log.Error("Failed to get folders", "err", err)
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -372,18 +389,14 @@ func (s *Service) GetMigrationDataJSON(ctx context.Context, id int64) ([]byte, e
|
|||||||
Data: f,
|
Data: f,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
migrationData := cloudmigration.MigrateDataRequestDTO{
|
migrationData := &cloudmigration.MigrateDataRequestDTO{
|
||||||
Items: migrationDataSlice,
|
Items: migrationDataSlice,
|
||||||
}
|
}
|
||||||
result, err := json.Marshal(migrationData)
|
|
||||||
if err != nil {
|
return migrationData, nil
|
||||||
s.log.Error("Failed to marshal datasources", "err", err)
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
return result, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) getDataSources(ctx context.Context, id int64) ([]datasources.AddDataSourceCommand, error) {
|
func (s *Service) getDataSources(ctx context.Context) ([]datasources.AddDataSourceCommand, error) {
|
||||||
dataSources, err := s.dsService.GetAllDataSources(ctx, &datasources.GetAllDataSourcesQuery{})
|
dataSources, err := s.dsService.GetAllDataSources(ctx, &datasources.GetAllDataSourcesQuery{})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
s.log.Error("Failed to get all datasources", "err", err)
|
s.log.Error("Failed to get all datasources", "err", err)
|
||||||
@@ -420,7 +433,7 @@ func (s *Service) getDataSources(ctx context.Context, id int64) ([]datasources.A
|
|||||||
return result, err
|
return result, err
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) getFolders(ctx context.Context, id int64) ([]folder.Folder, error) {
|
func (s *Service) getFolders(ctx context.Context) ([]folder.Folder, error) {
|
||||||
reqCtx := contexthandler.FromContext(ctx)
|
reqCtx := contexthandler.FromContext(ctx)
|
||||||
folders, err := s.folderService.GetFolders(ctx, folder.GetFoldersQuery{
|
folders, err := s.folderService.GetFolders(ctx, folder.GetFoldersQuery{
|
||||||
SignedInUser: reqCtx.SignedInUser,
|
SignedInUser: reqCtx.SignedInUser,
|
||||||
@@ -437,7 +450,7 @@ func (s *Service) getFolders(ctx context.Context, id int64) ([]folder.Folder, er
|
|||||||
return result, nil
|
return result, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) getDashboards(ctx context.Context, id int64) ([]dashboards.Dashboard, error) {
|
func (s *Service) getDashboards(ctx context.Context) ([]dashboards.Dashboard, error) {
|
||||||
dashs, err := s.dashboardService.GetAllDashboards(ctx)
|
dashs, err := s.dashboardService.GetAllDashboards(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -471,12 +484,26 @@ func (s *Service) GetMigrationStatus(ctx context.Context, id string, runID strin
|
|||||||
return cmr, nil
|
return cmr, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) GetMigrationStatusList(ctx context.Context, migrationID string) ([]*cloudmigration.CloudMigrationRun, error) {
|
func (s *Service) GetMigrationRunList(ctx context.Context, migrationID string) (*cloudmigration.CloudMigrationRunList, error) {
|
||||||
cmrs, err := s.store.GetMigrationStatusList(ctx, migrationID)
|
runs, err := s.store.GetMigrationStatusList(ctx, migrationID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("retrieving migration statuses from db: %w", err)
|
return nil, fmt.Errorf("retrieving migration statuses from db: %w", err)
|
||||||
}
|
}
|
||||||
return cmrs, nil
|
|
||||||
|
runList := &cloudmigration.CloudMigrationRunList{Runs: []cloudmigration.MigrateDataResponseDTO{}}
|
||||||
|
for _, s := range runs {
|
||||||
|
// attempt to bind the raw result to a list of response item DTOs
|
||||||
|
r := cloudmigration.MigrateDataResponseDTO{
|
||||||
|
Items: []cloudmigration.MigrateDataResponseItemDTO{},
|
||||||
|
}
|
||||||
|
if err := json.Unmarshal(s.Result, &r); err != nil {
|
||||||
|
return nil, fmt.Errorf("error unmarshalling migration response items: %w", err)
|
||||||
|
}
|
||||||
|
r.RunID = s.ID
|
||||||
|
runList.Runs = append(runList.Runs, r)
|
||||||
|
}
|
||||||
|
|
||||||
|
return runList, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) DeleteMigration(ctx context.Context, id int64) (*cloudmigration.CloudMigration, error) {
|
func (s *Service) DeleteMigration(ctx context.Context, id int64) (*cloudmigration.CloudMigration, error) {
|
||||||
@@ -487,7 +514,7 @@ func (s *Service) DeleteMigration(ctx context.Context, id int64) (*cloudmigratio
|
|||||||
return c, nil
|
return c, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *Service) ParseCloudMigrationConfig() (string, error) {
|
func (s *Service) parseCloudMigrationConfig() (string, error) {
|
||||||
if s.cfg == nil {
|
if s.cfg == nil {
|
||||||
return "", fmt.Errorf("cfg cannot be nil")
|
return "", fmt.Errorf("cfg cannot be nil")
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -42,7 +42,7 @@ func (s *NoopServiceImpl) GetMigrationStatus(ctx context.Context, id string, run
|
|||||||
return nil, cloudmigration.ErrFeatureDisabledError
|
return nil, cloudmigration.ErrFeatureDisabledError
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *NoopServiceImpl) GetMigrationStatusList(ctx context.Context, id string) ([]*cloudmigration.CloudMigrationRun, error) {
|
func (s *NoopServiceImpl) GetMigrationRunList(ctx context.Context, id string) (*cloudmigration.CloudMigrationRunList, error) {
|
||||||
return nil, cloudmigration.ErrFeatureDisabledError
|
return nil, cloudmigration.ErrFeatureDisabledError
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -54,10 +54,6 @@ func (s *NoopServiceImpl) SaveMigrationRun(ctx context.Context, cmr *cloudmigrat
|
|||||||
return -1, cloudmigration.ErrInternalNotImplementedError
|
return -1, cloudmigration.ErrInternalNotImplementedError
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *NoopServiceImpl) GetMigrationDataJSON(ctx context.Context, id int64) ([]byte, error) {
|
func (s *NoopServiceImpl) RunMigration(context.Context, int64) (*cloudmigration.MigrateDataResponseDTO, error) {
|
||||||
return nil, cloudmigration.ErrFeatureDisabledError
|
return nil, cloudmigration.ErrFeatureDisabledError
|
||||||
}
|
}
|
||||||
|
|
||||||
func (s *NoopServiceImpl) ParseCloudMigrationConfig() (string, error) {
|
|
||||||
return "", cloudmigration.ErrFeatureDisabledError
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -0,0 +1,114 @@
|
|||||||
|
package cmsclient
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"net/http"
|
||||||
|
|
||||||
|
"github.com/grafana/grafana/pkg/infra/log"
|
||||||
|
"github.com/grafana/grafana/pkg/services/cloudmigration"
|
||||||
|
)
|
||||||
|
|
||||||
|
type Client interface {
|
||||||
|
ValidateKey(context.Context, cloudmigration.CloudMigration) error
|
||||||
|
MigrateData(context.Context, cloudmigration.CloudMigration, cloudmigration.MigrateDataRequestDTO) (*cloudmigration.MigrateDataResponseDTO, error)
|
||||||
|
}
|
||||||
|
|
||||||
|
const logPrefix = "cloudmigration.cmsclient"
|
||||||
|
|
||||||
|
func NewCMSClient(domain string) Client {
|
||||||
|
return &clientImpl{
|
||||||
|
domain: domain,
|
||||||
|
log: log.New(logPrefix),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type clientImpl struct {
|
||||||
|
domain string
|
||||||
|
log *log.ConcreteLogger
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *clientImpl) ValidateKey(ctx context.Context, cm cloudmigration.CloudMigration) error {
|
||||||
|
logger := c.log.FromContext(ctx)
|
||||||
|
|
||||||
|
path := fmt.Sprintf("https://cms-%s.%s/cloud-migrations/api/v1/validate-key", cm.ClusterSlug, c.domain)
|
||||||
|
|
||||||
|
// validation is an empty POST to CMS with the authorization header included
|
||||||
|
req, err := http.NewRequest("POST", path, bytes.NewReader(nil))
|
||||||
|
if err != nil {
|
||||||
|
logger.Error("error creating http request for token validation", "err", err.Error())
|
||||||
|
return fmt.Errorf("http request error: %w", err)
|
||||||
|
}
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
req.Header.Set("Authorization", fmt.Sprintf("Bearer %d:%s", cm.StackID, cm.AuthToken))
|
||||||
|
|
||||||
|
client := &http.Client{}
|
||||||
|
resp, err := client.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
logger.Error("error sending http request for token validation", "err", err.Error())
|
||||||
|
return fmt.Errorf("http request error: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
defer func() {
|
||||||
|
if err := resp.Body.Close(); err != nil {
|
||||||
|
logger.Error("closing request body", "err", err.Error())
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
if resp.StatusCode != 200 {
|
||||||
|
var errResp map[string]any
|
||||||
|
if err := json.NewDecoder(resp.Body).Decode(&errResp); err != nil {
|
||||||
|
logger.Error("decoding error response", "err", err.Error())
|
||||||
|
} else {
|
||||||
|
return fmt.Errorf("token validation failure: %v", errResp)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *clientImpl) MigrateData(ctx context.Context, cm cloudmigration.CloudMigration, request cloudmigration.MigrateDataRequestDTO) (*cloudmigration.MigrateDataResponseDTO, error) {
|
||||||
|
logger := c.log.FromContext(ctx)
|
||||||
|
|
||||||
|
path := fmt.Sprintf("https://cms-%s.%s/cloud-migrations/api/v1/migrate-data", cm.ClusterSlug, c.domain)
|
||||||
|
|
||||||
|
body, err := json.Marshal(request)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("error marshaling request: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Send the request to cms with the associated auth token
|
||||||
|
req, err := http.NewRequest(http.MethodPost, path, bytes.NewReader(body))
|
||||||
|
if err != nil {
|
||||||
|
c.log.Error("error creating http request for cloud migration run", "err", err.Error())
|
||||||
|
return nil, fmt.Errorf("http request error: %w", err)
|
||||||
|
}
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
req.Header.Set("Authorization", fmt.Sprintf("Bearer %d:%s", cm.StackID, cm.AuthToken))
|
||||||
|
|
||||||
|
client := &http.Client{}
|
||||||
|
resp, err := client.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
c.log.Error("error sending http request for cloud migration run", "err", err.Error())
|
||||||
|
return nil, fmt.Errorf("http request error: %w", err)
|
||||||
|
} else if resp.StatusCode >= 400 {
|
||||||
|
c.log.Error("received error response for cloud migration run", "statusCode", resp.StatusCode)
|
||||||
|
return nil, fmt.Errorf("http request error: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
defer func() {
|
||||||
|
if err := resp.Body.Close(); err != nil {
|
||||||
|
logger.Error("closing request body: %w", err)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
var result cloudmigration.MigrateDataResponseDTO
|
||||||
|
if err := json.NewDecoder(resp.Body).Decode(&result); err != nil {
|
||||||
|
logger.Error("unmarshalling response body: %w", err)
|
||||||
|
return nil, fmt.Errorf("unmarshalling migration run response: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
return &result, nil
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user