Merge pull request #6961 from ioito/hotfix/qx-sku-sync

fix: 优化sku同步逻辑
This commit is contained in:
Zexi Li
2020-07-03 11:40:04 +08:00
committed by GitHub
5 changed files with 128 additions and 212 deletions
+1 -42
View File
@@ -480,7 +480,7 @@ func (manager *SDBInstanceSkuManager) SyncDBInstanceSkus(ctx context.Context, us
syncResult := compare.SyncResult{}
iskus, err := meta.GetDBInstanceSkusByRegion(region.ExternalId)
iskus, err := meta.GetDBInstanceSkusByRegionExternalId(region.ExternalId)
if err != nil {
syncResult.Error(err)
return syncResult
@@ -545,52 +545,11 @@ func (sku *SDBInstanceSku) syncWithCloudSku(ctx context.Context, userCred mcclie
return err
}
func (manager *SDBInstanceSkuManager) getZoneBySuffix(region *SCloudregion, suffix string) (*SZone, error) {
q := ZoneManager.Query().Equals("cloudregion_id", region.Id).Endswith("external_id", suffix)
count, err := q.CountWithError()
if err != nil {
return nil, err
}
if count > 1 {
return nil, fmt.Errorf("duplicate zone with suffix %s in region %s", suffix, region.Name)
}
if count == 0 {
return nil, fmt.Errorf("failed to found zone with suffix %s in region %s", suffix, region.Name)
}
zone := &SZone{}
return zone, q.First(zone)
}
func (manager *SDBInstanceSkuManager) newFromCloudSku(ctx context.Context, userCred mcclient.TokenCredential, isku SDBInstanceSku, region *SCloudregion) error {
sku := &isku
sku.SetModelManager(manager, sku)
sku.Id = "" //避免使用yunion meta的id,导致出现duplicate entry问题
sku.CloudregionId = region.Id
if len(isku.Zone1) > 0 {
zone, err := manager.getZoneBySuffix(region, isku.Zone1)
if err != nil {
return errors.Wrapf(err, "failed to get zone1 info by %s", isku.Zone1)
}
sku.Zone1 = zone.Id
}
if len(isku.Zone2) > 0 {
zone, err := manager.getZoneBySuffix(region, isku.Zone2)
if err != nil {
return errors.Wrapf(err, "failed to get zone1 info by %s", isku.Zone2)
}
sku.Zone2 = zone.Id
}
if len(isku.Zone3) > 0 {
zone, err := manager.getZoneBySuffix(region, isku.Zone3)
if err != nil {
return errors.Wrapf(err, "failed to get zone1 info by %s", isku.Zone3)
}
sku.Zone3 = zone.Id
}
return manager.TableSpec().Insert(ctx, sku)
}
+1 -2
View File
@@ -383,8 +383,7 @@ func (manager *SElasticcacheSkuManager) SyncElasticcacheSkus(ctx context.Context
syncResult := compare.SyncResult{}
extSkuMeta.SetRegionFilter(region)
extSkus, err := extSkuMeta.GetElasticCacheSkus()
extSkus, err := extSkuMeta.GetElasticCacheSkusByRegionExternalId(region.ExternalId)
if err != nil {
syncResult.Error(err)
return syncResult
+1 -1
View File
@@ -1237,7 +1237,7 @@ func (manager *SServerSkuManager) SyncServerSkus(ctx context.Context, userCred m
syncResult := compare.SyncResult{}
extSkus, err := extSkuMeta.GetServerSkus(region)
extSkus, err := extSkuMeta.GetServerSkusByRegionExternalId(region.ExternalId)
if err != nil {
syncResult.Error(err)
return syncResult
+123 -161
View File
@@ -25,6 +25,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
v "yunion.io/x/pkg/util/version"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/compute/options"
@@ -39,201 +40,162 @@ server: 虚拟机
elasticcache: 弹性缓存(redis&memcached)
*/
type SSkuResourcesMeta struct {
region *SCloudregion
caches map[string][]jsonutils.JSONObject
zoneCaches map[string]*SZone
regionCaches map[string]*SCloudregion
Server string
ElasticCache string
DBInstance string `json:"dbinstance"`
DBInstanceBase string `json:"dbinstance_base"`
DBInstanceBase string `json:"dbinstance_base"`
ServerBase string `json:"server_base"`
ElasticCacheBase string `json:"elastic_cache_base"`
}
func (self *SSkuResourcesMeta) GetServerSkus(region *SCloudregion) ([]SServerSku, error) {
self.SetRegionFilter(region)
func (self *SSkuResourcesMeta) getZoneIdBySuffix(zoneMaps map[string]string, suffix string) string {
for externalId, id := range zoneMaps {
if strings.HasSuffix(externalId, suffix) {
return id
}
}
return ""
}
result := []SServerSku{}
objs, err := self.get(self.Server)
func (self *SSkuResourcesMeta) GetDBInstanceSkusByRegionExternalId(regionExternalId string) ([]SDBInstanceSku, error) {
regionId, zoneMaps, err := self.GetRegionIdAndZoneMaps(regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "self.get")
return nil, errors.Wrap(err, "GetRegionIdAndZoneMaps")
}
for _, obj := range objs {
sku := SServerSku{}
err = obj.Unmarshal(&sku)
if err != nil {
return nil, errors.Wrap(err, "obj.Unmarshal")
}
// provider must not be empty
if len(sku.Provider) == 0 {
log.Debugf("source sku error: provider should not be empty. %#v", sku)
continue
}
// 处理数据
sku.Id = ""
r, err := self.fetchRegion(sku.CloudregionId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.GetServerSkus.fetchRegion")
}
sku.CloudregionId = r.GetId()
if len(sku.ZoneId) > 0 {
zone, err := self.fetchZone(sku.ZoneId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.GetServerSkus.fetchZone")
}
sku.ZoneId = zone.GetId()
}
result = append(result, sku)
}
return result, nil
}
func (self *SSkuResourcesMeta) GetDBInstanceSkusByRegion(regionId string) ([]SDBInstanceSku, error) {
result := []SDBInstanceSku{}
objs, err := self.getSkusByRegion(self.DBInstanceBase, regionId)
objs, err := self.getSkusByRegion(self.DBInstanceBase, regionExternalId)
if err != nil {
return nil, errors.Wrapf(err, "getSkusByRegion")
}
for _, obj := range objs {
sku := SDBInstanceSku{}
sku.SetModelManager(DBInstanceSkuManager, &sku)
err = obj.Unmarshal(&sku)
if err != nil {
return nil, errors.Wrapf(err, "obj.Unmarshal")
}
if len(sku.Zone1) > 0 {
zoneId := self.getZoneIdBySuffix(zoneMaps, sku.Zone1)
if len(zoneId) == 0 {
return nil, fmt.Errorf("invalid sku %s %s zone1: %s", sku.Id, sku.CloudregionId, sku.Zone1)
}
sku.Zone1 = zoneId
}
if len(sku.Zone2) > 0 {
zoneId := self.getZoneIdBySuffix(zoneMaps, sku.Zone2)
if len(zoneId) == 0 {
return nil, fmt.Errorf("invalid sku %s %s zone2: %s", sku.Id, sku.CloudregionId, sku.Zone2)
}
sku.Zone2 = zoneId
}
if len(sku.Zone3) > 0 {
zoneId := self.getZoneIdBySuffix(zoneMaps, sku.Zone3)
if len(zoneId) == 0 {
return nil, fmt.Errorf("invalid sku %s %s zone3: %s", sku.Id, sku.CloudregionId, sku.Zone3)
}
sku.Zone3 = zoneId
}
sku.Id = ""
sku.CloudregionId = regionId
result = append(result, sku)
}
return result, nil
}
func (self *SSkuResourcesMeta) GetElasticCacheSkus() ([]SElasticcacheSku, error) {
result := []SElasticcacheSku{}
objs, err := self.get(self.ElasticCache)
func (self *SSkuResourcesMeta) getCloudregion(regionExternalId string) (*SCloudregion, error) {
region, err := db.FetchByExternalId(CloudregionManager, regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "self.get(self.ElasticCache)")
return nil, errors.Wrapf(err, "db.FetchByExternalId(%s)", regionExternalId)
}
return region.(*SCloudregion), nil
}
func (self *SSkuResourcesMeta) GetRegionIdAndZoneMaps(regionExternalId string) (string, map[string]string, error) {
region, err := self.getCloudregion(regionExternalId)
if err != nil {
return "", nil, errors.Wrap(err, "getCloudregion")
}
zones, err := region.GetZones()
if err != nil {
return "", nil, errors.Wrap(err, "GetZones")
}
zoneMaps := map[string]string{}
for _, zone := range zones {
zoneMaps[zone.ExternalId] = zone.Id
}
return region.Id, zoneMaps, nil
}
func (self *SSkuResourcesMeta) GetServerSkusByRegionExternalId(regionExternalId string) ([]SServerSku, error) {
regionId, zoneMaps, err := self.GetRegionIdAndZoneMaps(regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "GetRegionIdAndZoneMaps")
}
result := []SServerSku{}
objs, err := self.getSkusByRegion(self.ServerBase, regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "getSkusByRegion")
}
for _, obj := range objs {
sku := SServerSku{}
sku.SetModelManager(ElasticcacheSkuManager, &sku)
err = obj.Unmarshal(&sku)
if err != nil {
return nil, errors.Wrapf(err, "obj.Unmarshal")
}
if len(sku.ZoneId) > 0 {
zoneId := self.getZoneIdBySuffix(zoneMaps, sku.ZoneId)
if len(zoneId) == 0 {
return nil, fmt.Errorf("invalid sku %s %s zoneId: %s", sku.Id, sku.CloudregionId, sku.ZoneId)
}
sku.ZoneId = zoneId
}
sku.Id = ""
sku.CloudregionId = regionId
result = append(result, sku)
}
return result, nil
}
func (self *SSkuResourcesMeta) GetElasticCacheSkusByRegionExternalId(regionExternalId string) ([]SElasticcacheSku, error) {
regionId, zoneMaps, err := self.GetRegionIdAndZoneMaps(regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "GetRegionIdAndZoneMaps")
}
result := []SElasticcacheSku{}
objs, err := self.getSkusByRegion(self.ElasticCacheBase, regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "getSkusByRegion")
}
for _, obj := range objs {
sku := SElasticcacheSku{}
sku.SetModelManager(ElasticcacheSkuManager, &sku)
err = obj.Unmarshal(&sku)
if err != nil {
return nil, errors.Wrap(err, "obj.Unmarshal")
return nil, errors.Wrapf(err, "obj.Unmarshal")
}
// 处理数据
sku.Id = ""
r, err := self.fetchRegion(sku.CloudregionId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.GetElasticCacheSkus.fetchRegion")
}
sku.CloudregionId = r.GetId()
if len(sku.ZoneId) > 0 {
zone, err := self.fetchZone(sku.ZoneId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.GetElasticCacheSkus.MasterZone")
zoneId := self.getZoneIdBySuffix(zoneMaps, sku.ZoneId)
if len(zoneId) == 0 {
return nil, fmt.Errorf("invalid sku %s %s master zoneId: %s", sku.Id, sku.CloudregionId, sku.ZoneId)
}
sku.ZoneId = zone.GetId()
sku.ZoneId = zoneId
}
if len(sku.SlaveZoneId) > 0 {
zone, err := self.fetchZone(sku.SlaveZoneId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.GetElasticCacheSkus.SlaveZone")
zoneId := self.getZoneIdBySuffix(zoneMaps, sku.SlaveZoneId)
if len(zoneId) == 0 {
return nil, fmt.Errorf("invalid sku %s %s slave zoneId: %s", sku.Id, sku.CloudregionId, sku.SlaveZoneId)
}
sku.SlaveZoneId = zone.GetId()
sku.ZoneId = zoneId
}
sku.Id = ""
sku.CloudregionId = regionId
result = append(result, sku)
}
return result, nil
}
func (self *SSkuResourcesMeta) fetchZone(zoneExternalId string) (*SZone, error) {
if self.zoneCaches == nil {
self.zoneCaches = map[string]*SZone{}
}
if z, ok := self.zoneCaches[zoneExternalId]; ok {
return z, nil
}
_zone, err := db.FetchByExternalId(ZoneManager, zoneExternalId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.fetchZone.FetchByExternalId")
}
z := _zone.(*SZone)
self.zoneCaches[zoneExternalId] = z
return z, nil
}
func (self *SSkuResourcesMeta) fetchRegion(regionExternalId string) (*SCloudregion, error) {
if self.regionCaches == nil {
self.regionCaches = map[string]*SCloudregion{}
}
if r, ok := self.regionCaches[regionExternalId]; ok {
return r, nil
}
_region, err := db.FetchByExternalId(CloudregionManager, regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.fetchRegion.FetchByExternalId")
}
r := _region.(*SCloudregion)
self.regionCaches[regionExternalId] = r
return r, nil
}
func (self *SSkuResourcesMeta) SetRegionFilter(region *SCloudregion) {
self.region = region
}
func (self *SSkuResourcesMeta) filterByRegion(items []jsonutils.JSONObject) []jsonutils.JSONObject {
if self.region == nil {
return items
}
ret := []jsonutils.JSONObject{}
for i := range items {
item := items[i]
regionId, _ := item.GetString("cloudregion_id")
if self.region.GetExternalId() != strings.TrimSpace(regionId) {
continue
}
ret = append(ret, item)
}
return ret
}
func (self *SSkuResourcesMeta) get(url string) ([]jsonutils.JSONObject, error) {
if self.caches == nil {
self.caches = map[string][]jsonutils.JSONObject{}
}
if items, ok := self.caches[url]; !ok || len(items) == 0 {
items, err := self._get(url)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.get")
}
self.caches[url] = items
}
items := self.caches[url]
return self.filterByRegion(items), nil
}
func (self *SSkuResourcesMeta) getSkusByRegion(base string, region string) ([]jsonutils.JSONObject, error) {
url := fmt.Sprintf("%s/%s.json", base, region)
items, err := self._get(url)
@@ -253,6 +215,9 @@ func (self *SSkuResourcesMeta) _get(url string) ([]jsonutils.JSONObject, error)
return nil, fmt.Errorf("SkuResourcesMeta.get.NewRequest %s", err)
}
userAgent := "vendor/yunion-OneCloud@" + v.Get().GitVersion
req.Header.Set("User-Agent", userAgent)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
@@ -273,7 +238,7 @@ func (self *SSkuResourcesMeta) _get(url string) ([]jsonutils.JSONObject, error)
var ret []jsonutils.JSONObject
err = jsonContent.Unmarshal(&ret)
if err != nil {
return nil, fmt.Errorf("SkuResourcesMeta.get.Unmarshal %s", err)
return nil, fmt.Errorf("SkuResourcesMeta.get.Unmarshal %s content: %s url: %s", err, jsonContent, url)
}
return ret, nil
@@ -312,7 +277,6 @@ func SyncElasticCacheSkus(ctx context.Context, userCred mcclient.TokenCredential
region := &cloudregions[i]
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)
@@ -336,7 +300,6 @@ func SyncElasticCacheSkusByRegion(ctx context.Context, userCred mcclient.TokenCr
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())
log.Infof(notes)
@@ -366,7 +329,6 @@ func SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, isSt
cloudregions := fetchSkuSyncCloudregions()
for i := range cloudregions {
region := &cloudregions[i]
meta.SetRegionFilter(region)
result := ServerSkuManager.SyncServerSkus(ctx, userCred, region, meta)
notes := fmt.Sprintf("SyncServerSkusByRegion %s result: %s", region.Name, result.Result())
log.Infof(notes)
@@ -105,12 +105,8 @@ func (self *CloudAccountSyncSkusTask) OnInit(ctx context.Context, obj db.IStanda
}
if syncFunc != nil {
if result := syncFunc(ctx, self.GetUserCred(), &region, meta); result.IsError() {
self.taskFailed(ctx, account, result.AllError())
return
} else {
log.Infof(result.Result())
}
result := syncFunc(ctx, self.GetUserCred(), &region, meta)
log.Infof("Sync %s %s skus for region %s result: %s", region.Provider, res, region.Name, result.Result())
}
}