From 5fd512cf119d03e724b2a331a80a18d2d9b67605 Mon Sep 17 00:00:00 2001 From: TangBin Date: Mon, 6 Jan 2020 19:38:32 +0800 Subject: [PATCH] cloud account support manually sync skus --- cmd/climc/shell/cloudaccounts.go | 30 +++++ pkg/cloudcommon/db/opslog.go | 1 + pkg/compute/models/cloudaccounts.go | 39 +++++++ pkg/compute/models/cloudsync.go | 6 +- pkg/compute/models/dbinstance_skus.go | 10 +- pkg/compute/models/elasticcache_skus.go | 31 +++++- pkg/compute/models/skus.go | 3 +- pkg/compute/models/skus_tools.go | 54 +++++---- .../tasks/cloud_account_sync_skus_task.go | 104 ++++++++++++++++++ 9 files changed, 247 insertions(+), 31 deletions(-) create mode 100644 pkg/compute/tasks/cloud_account_sync_skus_task.go diff --git a/cmd/climc/shell/cloudaccounts.go b/cmd/climc/shell/cloudaccounts.go index fb378a1902..5c8ff86a43 100644 --- a/cmd/climc/shell/cloudaccounts.go +++ b/cmd/climc/shell/cloudaccounts.go @@ -777,4 +777,34 @@ func init() { printObject(result) return nil }) + + type CloudaccountSyncSkusOptions struct { + ID string `help:"ID or Name of cloud account"` + RESOURCE string `help:"Resource of skus" choices:"serversku|elasticcachesku|dbinstance_sku"` + Force bool `help:"Force sync no matter what"` + Provider string `help:"provider to sync"` + Region string `help:"region to sync"` + } + R(&CloudaccountSyncSkusOptions{}, "cloud-account-sync-skus", "Sync skus of a cloud account", func(s *mcclient.ClientSession, args *CloudaccountSyncSkusOptions) error { + params := jsonutils.NewDict() + params.Set("resource", jsonutils.NewString(args.RESOURCE)) + if args.Force { + params.Add(jsonutils.JSONTrue, "force") + } + + if len(args.Provider) > 0 { + params.Add(jsonutils.NewString(args.Provider), "cloudprovider") + } + + if len(args.Region) > 0 { + params.Add(jsonutils.NewString(args.Region), "cloudregion") + } + + result, err := modules.Cloudaccounts.PerformAction(s, args.ID, "sync-skus", params) + if err != nil { + return err + } + printObject(result) + return nil + }) } diff --git a/pkg/cloudcommon/db/opslog.go b/pkg/cloudcommon/db/opslog.go index 9f67ce54fe..58d67efdd1 100644 --- a/pkg/cloudcommon/db/opslog.go +++ b/pkg/cloudcommon/db/opslog.go @@ -206,6 +206,7 @@ const ( ACT_SYNC_CLOUD_DISK = "sync_cloud_disk" ACT_SYNC_CLOUD_SERVER = "sync_cloud_server" + ACT_SYNC_CLOUD_SKUS = "sync_cloud_skus" ACT_SYNC_CLOUD_EIP = "sync_cloud_eip" ACT_SYNC_CLOUD_PROJECT = "sync_cloud_project" ACT_SYNC_CLOUD_ELASTIC_CACHE = "sync_cloud_elastic_cache" diff --git a/pkg/compute/models/cloudaccounts.go b/pkg/compute/models/cloudaccounts.go index ef169b2f34..b462703d81 100644 --- a/pkg/compute/models/cloudaccounts.go +++ b/pkg/compute/models/cloudaccounts.go @@ -45,6 +45,7 @@ import ( "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/auth" "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/util/choices" "yunion.io/x/onecloud/pkg/util/logclient" "yunion.io/x/onecloud/pkg/util/rbacutils" "yunion.io/x/onecloud/pkg/util/stringutils2" @@ -1748,3 +1749,41 @@ func guessBrandForHypervisor(hypervisor string) string { } return brands[0] } + +func (account *SCloudaccount) AllowPerformSyncSkus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return db.IsAllowPerform(rbacutils.ScopeSystem, userCred, account, "sync-skus") +} + +func (account *SCloudaccount) PerformSyncSkus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + if !account.Enabled { + return nil, httperrors.NewInvalidStatusError("Account disabled") + } + + dataDict := data.(*jsonutils.JSONDict) + resourceV := validators.NewStringChoicesValidator("resource", choices.NewChoices(ServerSkuManager.Keyword(), ElasticcacheSkuManager.Keyword(), DBInstanceSkuManager.Keyword())) + regionV := validators.NewModelIdOrNameValidator("cloudregion", "cloudregion", account.GetOwnerId()) + providerV := validators.NewModelIdOrNameValidator("cloudprovider", "cloudprovider", account.GetOwnerId()) + keyV := map[string]validators.IValidator{ + "resource": resourceV, + "cloudregion": regionV.Optional(true), + "cloudprovider": providerV.Optional(true), + } + + for _, v := range keyV { + if err := v.Validate(dataDict); err != nil { + return nil, err + } + } + + force, _ := data.Bool("force") + if account.CanSync() || force { + task, err := taskman.TaskManager.NewTask(ctx, "CloudAccountSyncSkusTask", account, userCred, dataDict, "", "", nil) + if err != nil { + return nil, errors.Wrapf(err, "CloudAccountSyncSkusTask") + } + + task.ScheduleRun(nil) + } + + return nil, nil +} diff --git a/pkg/compute/models/cloudsync.go b/pkg/compute/models/cloudsync.go index c4e073aa63..3297788d54 100644 --- a/pkg/compute/models/cloudsync.go +++ b/pkg/compute/models/cloudsync.go @@ -113,7 +113,7 @@ func syncRegionSkus(ctx context.Context, userCred mcclient.TokenCredential, loca if cnt == 0 { // 提前同步instance type.如果同步失败可能导致vm 内存显示为0 - if err = syncServerSkusByRegion(ctx, userCred, localRegion); err != nil { + if err = SyncServerSkusByRegion(ctx, userCred, localRegion); err != nil { msg := fmt.Sprintf("Get Skus for region %s failed %s", localRegion.GetName(), err) log.Errorln(msg) // 暂时不终止同步 @@ -135,7 +135,7 @@ func syncRegionSkus(ctx context.Context, userCred mcclient.TokenCredential, loca } if cnt == 0 { - syncElasticCacheSkusByRegion(ctx, userCred, localRegion) + SyncElasticCacheSkusByRegion(ctx, userCred, localRegion) } } } @@ -1003,7 +1003,7 @@ func syncPublicCloudProviderInfo( if !driver.GetFactory().NeedSyncSkuFromCloud() { syncRegionSkus(ctx, userCred, localRegion) - syncRegionDBInstanceSkus(ctx, userCred, localRegion.Id, true) + SyncRegionDBInstanceSkus(ctx, userCred, localRegion.Id, true) } else { syncSkusFromPrivateCloud(ctx, userCred, syncResults, provider, remoteRegion) } diff --git a/pkg/compute/models/dbinstance_skus.go b/pkg/compute/models/dbinstance_skus.go index 5bd49ea138..3f05603a04 100644 --- a/pkg/compute/models/dbinstance_skus.go +++ b/pkg/compute/models/dbinstance_skus.go @@ -331,7 +331,7 @@ func (manager *SDBInstanceSkuManager) GetDBInstanceSkus(provider, cloudregionId, return skus, nil } -func (manager *SDBInstanceSkuManager) syncDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, meta *SSkuResourcesMeta) compare.SyncResult { +func (manager *SDBInstanceSkuManager) SyncDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, meta *SSkuResourcesMeta) compare.SyncResult { lockman.LockClass(ctx, manager, db.GetLockClassKey(manager, userCred)) defer lockman.ReleaseClass(ctx, manager, db.GetLockClassKey(manager, userCred)) @@ -451,7 +451,7 @@ func (manager *SDBInstanceSkuManager) newFromCloudSku(ctx context.Context, userC return manager.TableSpec().Insert(sku) } -func syncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCredential, regionId string, isStart bool) { +func SyncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCredential, regionId string, isStart bool) { if isStart { q := DBInstanceSkuManager.Query() if len(regionId) > 0 { @@ -480,7 +480,7 @@ func syncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCreden return } - meta, err := fetchSkuResourcesMeta() + meta, err := FetchSkuResourcesMeta() if err != nil { log.Errorf("failed to fetch sku resource meta: %v", err) return @@ -491,7 +491,7 @@ func syncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCreden log.Infof("region %s(%s) not support dbinstance, skip sync", region.Name, region.Id) continue } - result := DBInstanceSkuManager.syncDBInstanceSkus(ctx, userCred, ®ion, meta) + result := DBInstanceSkuManager.SyncDBInstanceSkus(ctx, userCred, ®ion, meta) msg := result.Result() notes := fmt.Sprintf("SyncDBInstanceSkus for region %s result: %s", region.Name, msg) log.Infof(notes) @@ -500,5 +500,5 @@ func syncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCreden } func SyncDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { - syncRegionDBInstanceSkus(ctx, userCred, "", isStart) + SyncRegionDBInstanceSkus(ctx, userCred, "", isStart) } diff --git a/pkg/compute/models/elasticcache_skus.go b/pkg/compute/models/elasticcache_skus.go index 7fb92adef8..a1cf776764 100644 --- a/pkg/compute/models/elasticcache_skus.go +++ b/pkg/compute/models/elasticcache_skus.go @@ -229,7 +229,7 @@ func (manager *SElasticcacheSkuManager) FetchSkusByRegion(regionID string) ([]SE return skus, nil } -func (manager *SElasticcacheSkuManager) syncElasticcacheSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSkuMeta *SSkuResourcesMeta) compare.SyncResult { +func (manager *SElasticcacheSkuManager) SyncElasticcacheSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSkuMeta *SSkuResourcesMeta) compare.SyncResult { lockman.LockClass(ctx, manager, db.GetLockClassKey(manager, userCred)) defer lockman.ReleaseClass(ctx, manager, db.GetLockClassKey(manager, userCred)) @@ -434,3 +434,32 @@ func (manager *SElasticcacheSkuManager) GetPropertyCapability(ctx context.Contex return result, nil } + +func (manager *SElasticcacheSkuManager) PerformActionSync(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) { + data := query.(*jsonutils.JSONDict) + cloudprovider := validators.NewModelIdOrNameValidator("cloudprovider", "cloudprovider", nil) + cloudregion := validators.NewModelIdOrNameValidator("cloudregion", "cloudregion", nil) + + keyV := map[string]validators.IValidator{ + "provider": cloudprovider.Optional(true), + "cloudregion": cloudregion.Optional(true), + } + + for _, v := range keyV { + if err := v.Validate(data); err != nil { + return nil, err + } + } + + regions := []SCloudregion{} + if region, err := data.GetString("cloudregion"); err == nil && len(region) > 0 { + regions = append(regions, *cloudregion.Model.(*SCloudregion)) + } else if provider, err := data.GetString("cloudprovider"); err == nil && len(provider) > 0 { + regions, err = CloudregionManager.GetRegionByProvider(provider) + if err != nil { + return nil, err + } + } + + return nil, nil +} diff --git a/pkg/compute/models/skus.go b/pkg/compute/models/skus.go index 8eb7a4304e..d66bc184ce 100644 --- a/pkg/compute/models/skus.go +++ b/pkg/compute/models/skus.go @@ -1125,6 +1125,7 @@ func (self *SServerSku) PerformDisable(ctx context.Context, userCred mcclient.To func (self *SServerSku) syncWithCloudSku(ctx context.Context, userCred mcclient.TokenCredential, extSku SServerSku) error { _, err := db.Update(self, func() error { + self.ZoneId = extSku.ZoneId self.PrepaidStatus = extSku.PrepaidStatus self.PostpaidStatus = extSku.PostpaidStatus return nil @@ -1155,7 +1156,7 @@ func (manager *SServerSkuManager) FetchSkusByRegion(regionID string) ([]SServerS return skus, nil } -func (manager *SServerSkuManager) syncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSkuMeta *SSkuResourcesMeta) compare.SyncResult { +func (manager *SServerSkuManager) SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSkuMeta *SSkuResourcesMeta) compare.SyncResult { lockman.LockClass(ctx, manager, db.GetLockClassKey(manager, userCred)) defer lockman.ReleaseClass(ctx, manager, db.GetLockClassKey(manager, userCred)) diff --git a/pkg/compute/models/skus_tools.go b/pkg/compute/models/skus_tools.go index fcbfd64600..25f2dc6c38 100644 --- a/pkg/compute/models/skus_tools.go +++ b/pkg/compute/models/skus_tools.go @@ -318,34 +318,46 @@ func SyncElasticCacheSkus(ctx context.Context, userCred mcclient.TokenCredential } } - meta, err := fetchSkuResourcesMeta() + meta, err := FetchSkuResourcesMeta() if err != nil { - log.Errorf("SyncElasticCacheSkus.fetchSkuResourcesMeta %s", err) + log.Errorf("SyncElasticCacheSkus.FetchSkuResourcesMeta %s", err) return } cloudregions := fetchSkuSyncCloudregions() for i := range cloudregions { region := &cloudregions[i] - meta.SetRegionFilter(region) - result := ElasticcacheSkuManager.syncElasticcacheSkus(ctx, userCred, region, meta) - notes := fmt.Sprintf("syncElasticCacheSkusByRegion %s result: %s", region.Name, result.Result()) - log.Infof(notes) + + if region.GetDriver().IsSupportedElasticcache() { + meta.SetRegionFilter(region) + result := ElasticcacheSkuManager.SyncElasticcacheSkus(ctx, userCred, region, meta) + notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s result: %s", region.Name, result.Result()) + log.Infof(notes) + } else { + notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s not support elasticcache", region.Name) + log.Infof(notes) + } } } // 同步Region elasticcache sku列表. -func syncElasticCacheSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion) { - meta, err := fetchSkuResourcesMeta() +func SyncElasticCacheSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion) error { + if region.GetDriver().IsSupportedElasticcache() { + notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s not support elasticcache", region.Name) + log.Infof(notes) + return nil + } + + meta, err := FetchSkuResourcesMeta() if err != nil { - log.Errorf("syncElasticCacheSkusByRegion.fetchSkuResourcesMeta %s", err) - return + return errors.Wrap(err, "SyncElasticCacheSkusByRegion.FetchSkuResourcesMeta") } meta.SetRegionFilter(region) - result := ElasticcacheSkuManager.syncElasticcacheSkus(ctx, userCred, region, meta) - notes := fmt.Sprintf("syncElasticCacheSkusByRegion %s result: %s", region.Name, result.Result()) + result := ElasticcacheSkuManager.SyncElasticcacheSkus(ctx, userCred, region, meta) + notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s result: %s", region.Name, result.Result()) log.Infof(notes) + return nil } // 全量同步sku列表. @@ -362,9 +374,9 @@ func SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, isSt } } - meta, err := fetchSkuResourcesMeta() + meta, err := FetchSkuResourcesMeta() if err != nil { - log.Errorf("SyncServerSkus.fetchSkuResourcesMeta %s", err) + log.Errorf("SyncServerSkus.FetchSkuResourcesMeta %s", err) return } @@ -372,7 +384,7 @@ func SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, isSt for i := range cloudregions { region := &cloudregions[i] meta.SetRegionFilter(region) - result := ServerSkuManager.syncServerSkus(ctx, userCred, region, meta) + result := ServerSkuManager.SyncServerSkus(ctx, userCred, region, meta) notes := fmt.Sprintf("SyncServerSkusByRegion %s result: %s", region.Name, result.Result()) log.Infof(notes) } @@ -383,19 +395,19 @@ func SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, isSt } // 同步指定region sku列表 -func syncServerSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion) error { - meta, err := fetchSkuResourcesMeta() +func SyncServerSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion) error { + meta, err := FetchSkuResourcesMeta() if err != nil { - return errors.Wrap(err, "syncServerSkusByRegion.fetchSkuResourcesMeta") + return errors.Wrap(err, "SyncServerSkusByRegion.FetchSkuResourcesMeta") } - result := ServerSkuManager.syncServerSkus(ctx, userCred, region, meta) - notes := fmt.Sprintf("syncServerSkusByRegion %s result: %s", region.Name, result.Result()) + result := ServerSkuManager.SyncServerSkus(ctx, userCred, region, meta) + notes := fmt.Sprintf("SyncServerSkusByRegion %s result: %s", region.Name, result.Result()) log.Infof(notes) return nil } -func fetchSkuResourcesMeta() (*SSkuResourcesMeta, error) { +func FetchSkuResourcesMeta() (*SSkuResourcesMeta, error) { s := auth.GetAdminSession(context.Background(), options.Options.Region, "") meta, err := modules.OfflineCloudmeta.GetSkuSourcesMeta(s) if err != nil { diff --git a/pkg/compute/tasks/cloud_account_sync_skus_task.go b/pkg/compute/tasks/cloud_account_sync_skus_task.go new file mode 100644 index 0000000000..fc7bd0dd26 --- /dev/null +++ b/pkg/compute/tasks/cloud_account_sync_skus_task.go @@ -0,0 +1,104 @@ +package tasks + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/util/compare" + "yunion.io/x/pkg/utils" + + api "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/util/logclient" +) + +type CloudAccountSyncSkusTask struct { + taskman.STask +} + +func init() { + taskman.RegisterTask(CloudAccountSyncSkusTask{}) +} + +func (self *CloudAccountSyncSkusTask) taskFailed(ctx context.Context, account *models.SCloudaccount, err error) { + account.SetStatus(self.UserCred, api.CLOUD_PROVIDER_SYNC_STATUS_ERROR, err.Error()) + db.OpsLog.LogEvent(account, db.ACT_SYNC_CLOUD_SKUS, err.Error(), self.GetUserCred()) + logclient.AddActionLogWithStartable(self, account, logclient.ACT_CLOUD_SYNC, err.Error(), self.UserCred, false) + self.SetStageFailed(ctx, err.Error()) +} + +func (self *CloudAccountSyncSkusTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + account := obj.(*models.SCloudaccount) + + regions := []models.SCloudregion{} + if regionId, _ := self.GetParams().GetString("cloudregion_id"); len(regionId) > 0 { + _region, err := db.FetchById(models.CloudregionManager, regionId) + if err != nil { + self.taskFailed(ctx, account, err) + return + } + + region := _region.(*models.SCloudregion) + regions = append(regions, *region) + } else if providerId, _ := self.GetParams().GetString("cloudprovider_id"); len(providerId) > 0 { + provider, err := db.FetchById(models.CloudproviderManager, providerId) + if err != nil { + self.taskFailed(ctx, account, err) + return + } + + _regions := provider.(*models.SCloudprovider).GetCloudproviderRegions() + for i := range _regions { + region := _regions[i].GetRegion() + regions = append(regions, *region) + } + } else { + providers := account.GetEnabledCloudproviders() + for _, provider := range providers { + ids := []string{} + _regions := provider.GetCloudproviderRegions() + for i := range _regions { + region := _regions[i].GetRegion() + if region != nil && !utils.IsInStringArray(region.GetId(), ids) { + regions = append(regions, *region) + ids = append(ids, region.GetId()) + } + } + } + } + + res, _ := self.GetParams().GetString("resource") + meta, err := models.FetchSkuResourcesMeta() + if err != nil { + self.taskFailed(ctx, account, err) + return + } + + type SyncFunc func(ctx context.Context, userCred mcclient.TokenCredential, region *models.SCloudregion, extSkuMeta *models.SSkuResourcesMeta) compare.SyncResult + var syncFunc SyncFunc + for _, region := range regions { + switch res { + case models.ServerSkuManager.Keyword(): + syncFunc = models.ServerSkuManager.SyncServerSkus + case models.ElasticcacheSkuManager.Keyword(): + syncFunc = models.ElasticcacheSkuManager.SyncElasticcacheSkus + case models.ElasticcacheSkuManager.Keyword(): + syncFunc = models.DBInstanceSkuManager.SyncDBInstanceSkus + } + + if syncFunc != nil { + if result := syncFunc(ctx, self.GetUserCred(), ®ion, meta); result.IsError() { + self.taskFailed(ctx, account, result.AllError()) + return + } else { + log.Infof(result.Result()) + } + } + } + + self.SetStageComplete(ctx, nil) +}