fix(region): optimized sku sync

This commit is contained in:
ioito
2023-06-12 18:04:49 +08:00
parent e6a4361c92
commit dc388928e7
18 changed files with 791 additions and 1017 deletions
+13 -27
View File
@@ -47,6 +47,7 @@ import (
modules "yunion.io/x/onecloud/pkg/mcclient/modules/compute"
"yunion.io/x/onecloud/pkg/util/logclient"
"yunion.io/x/onecloud/pkg/util/stringutils2"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
// +onecloud:swagger-gen-ignore
@@ -850,10 +851,6 @@ func (self *SCloudaccount) GetCloudaccountByProvider(provider string) ([]SClouda
return accounts, nil
}
const (
EMPTY_MD5 = "d751713988987e9331980363e24189ce"
)
func (self *SCloudaccount) SyncSystemCloudpoliciesFromCloud(ctx context.Context, userCred mcclient.TokenCredential, refresh bool) error {
dbPolicies, err := self.GetSystemCloudpolicies()
if err != nil {
@@ -864,13 +861,14 @@ func (self *SCloudaccount) SyncSystemCloudpoliciesFromCloud(ctx context.Context,
return nil
}
s := auth.GetAdminSession(context.Background(), options.Options.Region)
transport := httputils.GetTransport(true)
transport.Proxy = options.Options.HttpTransportProxyFunc()
client := &http.Client{Transport: transport}
policyBase, index, err := modules.OfflineCloudmeta.GetSkuIndex(s, client, "cloudpolicy_base")
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return errors.Wrapf(err, "get cloudpolicy")
return errors.Wrapf(err, "FetchYunionmeta")
}
index, err := meta.Index(CloudpolicyManager.Keyword())
if err != nil {
return errors.Wrapf(err, "Index")
}
skuMeta := &SCloudpolicy{}
@@ -879,28 +877,16 @@ func (self *SCloudaccount) SyncSystemCloudpoliciesFromCloud(ctx context.Context,
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[self.Provider]
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
log.Debugf("%s cloudpolicy is empty skip syncing", self.Provider)
return nil
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Debugf("%s cloudpolicy not changed skip syncing", self.Provider)
if !ok || newMd5 == yunionmeta.EMPTY_MD5 || len(oldMd5) > 0 && newMd5 == oldMd5 {
return nil
}
policyUrl := strings.TrimSuffix(policyBase, "/") + fmt.Sprintf("/%s.json", self.Provider)
_, body, err := httputils.JSONRequest(client, ctx, httputils.GET, policyUrl, nil, nil, false)
if err != nil {
return errors.Wrapf(err, "JSONRequest(%s)", policyUrl)
}
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
iPolicies := []SCloudpolicy{}
err = body.Unmarshal(&iPolicies)
err = meta.List(CloudpolicyManager.Keyword(), self.Provider, &iPolicies)
if err != nil {
return errors.Wrapf(err, "body.Unmarshal")
return errors.Wrapf(err, "meta.List")
}
return self.syncSystemCloudpoliciesFromCloud(ctx, userCred, iPolicies, dbPolicies)
+21 -17
View File
@@ -17,12 +17,14 @@ package models
import (
"context"
"database/sql"
"fmt"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
type SCloudimageManager struct {
@@ -69,13 +71,13 @@ func SyncPublicCloudImages(ctx context.Context, userCred mcclient.TokenCredentia
return
}
meta, err := FetchSkuResourcesMeta()
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
log.Errorf("SyncServerSkus.FetchSkuResourcesMeta %s", err)
log.Errorf("FetchYunionmeta %v", err)
return
}
index, err := meta.getServerSkuIndex()
index, err := meta.Index(CloudimageManager.Keyword())
if err != nil {
log.Errorf("getServerSkuIndex error: %v", err)
return
@@ -90,25 +92,17 @@ func SyncPublicCloudImages(ctx context.Context, userCred mcclient.TokenCredentia
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
log.Infof("%s images is empty skip syncing", region.Name)
if !ok || newMd5 == yunionmeta.EMPTY_MD5 || len(oldMd5) > 0 && newMd5 == oldMd5 {
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Infof("%s cloud images not changed skip syncing", region.Name)
continue
}
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
err = regions[i].SyncCloudImages(ctx, userCred, !isStart, false)
if err != nil {
log.Errorf("SyncCloudImages for region %s(%s) error: %v", regions[i].Name, regions[i].Id, err)
continue
}
storagecaches, err := regions[i].GetStoragecaches()
if err != nil {
log.Errorf("GetStoragecaches for region %s(%s) error: %v", regions[i].Name, regions[i].Id, err)
@@ -141,17 +135,27 @@ func (self *SCloudimage) syncRemove(ctx context.Context, userCred mcclient.Token
return self.Delete(ctx, userCred)
}
func (self *SCloudimage) syncWithImage(ctx context.Context, userCred mcclient.TokenCredential, image SCachedimage) error {
func (self *SCloudimage) syncWithImage(ctx context.Context, userCred mcclient.TokenCredential, image SCachedimage, region *SCloudregion) error {
_cachedImage, err := db.FetchByExternalId(CachedimageManager, image.GetGlobalId())
if err != nil {
if errors.Cause(err) != sql.ErrNoRows {
return errors.Wrapf(err, "db.FetchByExternalId(%s)", image.GetGlobalId())
}
image := &image
image.Id = ""
image.SetModelManager(CachedimageManager, image)
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return err
}
skuUrl := fmt.Sprintf("%s/%s/%s.json", meta.ImageBase, region.ExternalId, image.GetGlobalId())
err = meta.Get(skuUrl, image)
if err != nil {
return errors.Wrapf(err, "Get")
}
image.IsPublic = true
image.ProjectId = "system"
image.SetModelManager(CachedimageManager, image)
err = CachedimageManager.TableSpec().Insert(ctx, image)
if err != nil {
return errors.Wrapf(err, "Insert cachedimage")
+19 -6
View File
@@ -17,6 +17,7 @@ package models
import (
"context"
"database/sql"
"fmt"
"strings"
"time"
@@ -38,6 +39,7 @@ import (
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/util/stringutils2"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
type SCloudregionManager struct {
@@ -1159,11 +1161,12 @@ func (self *SCloudregion) SyncCloudImages(ctx context.Context, userCred mcclient
if len(dbImages) > 0 && systemImageCount > 0 && !refresh {
return nil
}
meta, err := FetchSkuResourcesMeta()
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return errors.Wrapf(err, "FetchSkuResourcesMeta")
return errors.Wrapf(err, "FetchYunionmeta")
}
iImages, err := meta.GetCloudimages(self.ExternalId)
iImages := []SCachedimage{}
err = meta.List(CloudimageManager.Keyword(), self.ExternalId, &iImages)
if err != nil {
return errors.Wrapf(err, "GetCloudimages")
}
@@ -1190,7 +1193,7 @@ func (self *SCloudregion) SyncCloudImages(ctx context.Context, userCred mcclient
if !xor {
for i := 0; i < len(commonext); i++ {
err := commondb[i].syncWithImage(ctx, userCred, commonext[i])
err := commondb[i].syncWithImage(ctx, userCred, commonext[i], self)
if err != nil {
result.UpdateError(errors.Wrapf(err, "updateCachedImage"))
continue
@@ -1230,10 +1233,20 @@ func (self *SCloudregion) newCloudimage(ctx context.Context, userCred mcclient.T
return errors.Wrapf(err, "db.FetchModelObjects(%s)", iImage.GetGlobalId())
}
image := &iImage
image.Id = ""
image.SetModelManager(CachedimageManager, image)
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return err
}
skuUrl := fmt.Sprintf("%s/%s/%s.json", meta.ImageBase, self.ExternalId, iImage.GetGlobalId())
err = meta.Get(skuUrl, image)
if err != nil {
return errors.Wrapf(err, "Get")
}
image.IsPublic = true
image.ProjectId = "system"
image.SetModelManager(CachedimageManager, image)
err = CachedimageManager.TableSpec().Insert(ctx, image)
if err != nil {
return errors.Wrapf(err, "Insert cachedimage")
+1 -1
View File
@@ -162,7 +162,7 @@ func syncRegionSkus(ctx context.Context, userCred mcclient.TokenCredential, loca
if cnt == 0 {
// 提前同步instance type.如果同步失败可能导致vm 内存显示为0
if ret := SyncServerSkusByRegion(ctx, userCred, localRegion, nil, xor); ret.IsError() {
if ret := SyncServerSkusByRegion(ctx, userCred, localRegion, xor); ret.IsError() {
msg := fmt.Sprintf("Get Skus for region %s failed %s", localRegion.GetName(), ret.Result())
log.Errorln(msg)
// 暂时不终止同步
+81 -41
View File
@@ -35,6 +35,7 @@ import (
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/util/stringutils2"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
type SDBInstanceSkuManager struct {
@@ -484,24 +485,30 @@ func (manager *SDBInstanceSkuManager) SyncDBInstanceSkus(
ctx context.Context,
userCred mcclient.TokenCredential,
region *SCloudregion,
meta *SSkuResourcesMeta,
xor bool,
) compare.SyncResult {
lockman.LockRawObject(ctx, manager.Keyword(), region.Id)
defer lockman.ReleaseRawObject(ctx, manager.Keyword(), region.Id)
syncResult := compare.SyncResult{}
result := compare.SyncResult{}
iskus, err := meta.GetDBInstanceSkusByRegionExternalId(region.ExternalId)
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(errors.Wrapf(err, "FetchYunionmeta"))
return result
}
iskus := []SDBInstanceSku{}
err = meta.List(manager.Keyword(), region.ExternalId, &iskus)
if err != nil {
result.Error(err)
return result
}
dbSkus, err := manager.fetchDBInstanceSkus(region.Provider, region)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(err)
return result
}
removed := make([]SDBInstanceSku, 0)
@@ -511,37 +518,37 @@ func (manager *SDBInstanceSkuManager) SyncDBInstanceSkus(
err = compare.CompareSets(dbSkus, iskus, &removed, &commondb, &commonext, &added)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(err)
return result
}
for i := 0; i < len(removed); i += 1 {
err = removed[i].Delete(ctx, userCred)
if err != nil {
syncResult.DeleteError(err)
result.DeleteError(err)
} else {
syncResult.Delete()
result.Delete()
}
}
if !xor {
for i := 0; i < len(commondb); i += 1 {
err = commondb[i].syncWithCloudSku(ctx, userCred, commonext[i])
if err != nil {
syncResult.UpdateError(err)
result.UpdateError(err)
} else {
syncResult.Update()
result.Update()
}
}
}
for i := 0; i < len(added); i += 1 {
err = manager.newFromCloudSku(ctx, userCred, added[i], region)
err = region.newDBInstanceSkuFromCloudSku(ctx, userCred, added[i].GetGlobalId())
if err != nil {
syncResult.AddError(err)
result.AddError(err)
} else {
syncResult.Add()
result.Add()
}
}
return syncResult
return result
}
func (sku SDBInstanceSku) GetGlobalId() string {
@@ -551,21 +558,62 @@ func (sku SDBInstanceSku) GetGlobalId() string {
func (sku *SDBInstanceSku) syncWithCloudSku(ctx context.Context, userCred mcclient.TokenCredential, isku SDBInstanceSku) error {
_, err := db.Update(sku, func() error {
sku.Status = isku.Status
sku.TPS = isku.TPS
sku.QPS = isku.QPS
sku.MultiAZ = isku.MultiAZ
sku.MaxConnections = isku.MaxConnections
return nil
})
return err
}
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
return manager.TableSpec().Insert(ctx, sku)
func (self *SCloudregion) newDBInstanceSkuFromCloudSku(ctx context.Context, userCred mcclient.TokenCredential, externalId string) error {
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return err
}
zones, err := self.GetZones()
if err != nil {
return errors.Wrap(err, "GetZones")
}
zoneMaps := map[string]string{}
for _, zone := range zones {
zoneMaps[zone.ExternalId] = zone.Id
}
sku := &SDBInstanceSku{}
sku.SetModelManager(DBInstanceSkuManager, sku)
skuUrl := fmt.Sprintf("%s/%s/%s.json", meta.DBInstanceBase, self.ExternalId, externalId)
err = meta.Get(skuUrl, sku)
if err != nil {
return errors.Wrapf(err, "Get")
}
if len(sku.Zone1) > 0 {
zoneId := yunionmeta.GetZoneIdBySuffix(zoneMaps, sku.Zone1)
if len(zoneId) == 0 {
return errors.Wrapf(err, "GetZoneIdBySuffix(%s)", sku.Zone1)
}
sku.Zone1 = zoneId
}
if len(sku.Zone2) > 0 {
zoneId := yunionmeta.GetZoneIdBySuffix(zoneMaps, sku.Zone2)
if len(zoneId) == 0 {
return errors.Wrapf(err, "GetZoneIdBySuffix(%s)", sku.Zone2)
}
sku.Zone2 = zoneId
}
if len(sku.Zone3) > 0 {
zoneId := yunionmeta.GetZoneIdBySuffix(zoneMaps, sku.Zone3)
if len(zoneId) == 0 {
return errors.Wrapf(err, "GetZoneIdBySuffix(%s)", sku.Zone3)
}
sku.Zone3 = zoneId
}
sku.CloudregionId = self.Id
sku.Provider = self.Provider
sku.Enabled = tristate.True
return DBInstanceSkuManager.TableSpec().Insert(ctx, sku)
}
func SyncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCredential, regionId string, isStart, xor bool) {
@@ -597,13 +645,13 @@ func SyncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCreden
return
}
meta, err := FetchSkuResourcesMeta()
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
log.Errorf("failed to fetch sku resource meta: %v", err)
log.Errorf("FetchYunionmeta: %v", err)
return
}
index, err := meta.getSkuIndex(meta.DBInstanceBase)
index, err := meta.Index(DBInstanceSkuManager.Keyword())
if err != nil {
log.Errorf("get rds sku index error: %v", err)
return
@@ -620,21 +668,13 @@ func SyncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCreden
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
log.Debugf("%s DBInstance skus is empty skip syncing", region.Name)
if !ok || newMd5 == yunionmeta.EMPTY_MD5 || len(oldMd5) > 0 && newMd5 == oldMd5 {
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Debugf("%s DBInstance skus not changed skip syncing", region.Name)
continue
}
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
result := DBInstanceSkuManager.SyncDBInstanceSkus(ctx, userCred, &region, meta, xor)
result := DBInstanceSkuManager.SyncDBInstanceSkus(ctx, userCred, &region, xor)
msg := result.Result()
notes := fmt.Sprintf("sync rds sku for region %s result: %s", region.Name, msg)
log.Debugf(notes)
@@ -1011,9 +1011,6 @@ func (self *SElasticcache) ValidatorChangeSpecData(ctx context.Context, userCred
}
sku := skuV.Model.(*SElasticcacheSku)
if sku.Provider != self.GetProviderName() {
return nil, httperrors.NewInputParameterError("provider mismatch: %s instance can't use %s sku", self.GetProviderName(), sku.Provider)
}
region, _ := self.GetRegion()
if sku.CloudregionId != region.Id {
+136 -70
View File
@@ -17,6 +17,7 @@ package models
import (
"context"
"database/sql"
"fmt"
"strings"
"yunion.io/x/cloudmux/pkg/cloudprovider"
@@ -36,6 +37,7 @@ import (
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/util/stringutils2"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
type SElasticcacheSkuManager struct {
@@ -186,67 +188,6 @@ func (manager *SElasticcacheSkuManager) GetSkuCountByRegion(regionId string) (in
return q.CountWithError()
}
/*func (manager *SElasticcacheSkuManager) FetchCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, objs []db.IModel, fields stringutils2.SSortedStrings) []*jsonutils.JSONDict {
regions := map[string]string{}
for i := range objs {
cloudregionId := objs[i].(*SElasticcacheSku).CloudregionId
if _, ok := regions[cloudregionId]; !ok {
regions[cloudregionId] = cloudregionId
}
}
regionIds := []string{}
for k, _ := range regions {
regionIds = append(regionIds, regions[k])
}
if len(regionIds) == 0 {
return nil
}
regionObjs := []SCloudregion{}
err := CloudregionManager.Query().In("id", regionIds).All(&regionObjs)
if err != nil {
log.Errorf("elasticcacheSkuManager.FetchCustomizeColumns %s", err)
return nil
}
for i := range regionObjs {
regionObj := regionObjs[i]
regions[regionObj.Id] = regionObj.Name
}
ret := []*jsonutils.JSONDict{}
for i := range objs {
cloudregionId := objs[i].(*SElasticcacheSku).CloudregionId
fileds := jsonutils.NewDict()
fileds.Set("region", jsonutils.NewString(regions[cloudregionId]))
if region, err := db.FetchById(CloudregionManager, cloudregionId); err == nil {
fileds.Set("region_external_id", jsonutils.NewString(region.(*SCloudregion).ExternalId))
segs := strings.Split(region.(*SCloudregion).ExternalId, "/")
if len(segs) >= 2 {
fileds.Set("region_ext_id", jsonutils.NewString(segs[1]))
}
}
zoneId := objs[i].(*SElasticcacheSku).ZoneId
if len(zoneId) > 0 {
if zone, err := db.FetchById(ZoneManager, zoneId); err == nil {
fileds.Set("zone_external_id", jsonutils.NewString(zone.(*SZone).ExternalId))
segs := strings.Split(zone.(*SZone).ExternalId, "/")
if len(segs) >= 3 {
fileds.Set("zone_ext_id", jsonutils.NewString(segs[2]))
}
}
}
ret = append(ret, fileds)
}
return ret
}*/
// 弹性缓存套餐规格列表
func (manager *SElasticcacheSkuManager) ListItemFilter(
ctx context.Context,
@@ -413,13 +354,24 @@ func (manager *SElasticcacheSkuManager) FetchSkusByRegion(regionID string) ([]SE
return skus, nil
}
func (manager *SElasticcacheSkuManager) SyncElasticcacheSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSkuMeta *SSkuResourcesMeta, xor bool) compare.SyncResult {
func (self *SElasticcacheSku) GetElasticcacheCount() (int, error) {
q := ElasticcacheManager.Query().Equals("instance_type", self.Name).Equals("zone_id", self.ZoneId)
return q.CountWithError()
}
func (manager *SElasticcacheSkuManager) SyncElasticcacheSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, xor bool) compare.SyncResult {
lockman.LockRawObject(ctx, manager.Keyword(), region.Id)
defer lockman.ReleaseRawObject(ctx, manager.Keyword(), region.Id)
syncResult := compare.SyncResult{}
extSkus, err := extSkuMeta.GetElasticCacheSkusByRegionExternalId(region.ExternalId)
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return syncResult
}
extSkus := []SElasticcacheSku{}
err = meta.List(manager.Keyword(), region.ExternalId, &extSkus)
if err != nil {
syncResult.Error(err)
return syncResult
@@ -443,7 +395,11 @@ func (manager *SElasticcacheSkuManager) SyncElasticcacheSkus(ctx context.Context
}
for i := 0; i < len(removed); i += 1 {
err = removed[i].MarkAsSoldout(ctx)
if cnt, _ := removed[i].GetElasticcacheCount(); cnt > 0 {
err = removed[i].MarkAsSoldout(ctx)
} else {
err = removed[i].Delete(ctx, userCred)
}
if err != nil {
syncResult.DeleteError(err)
} else {
@@ -461,7 +417,7 @@ func (manager *SElasticcacheSkuManager) SyncElasticcacheSkus(ctx context.Context
}
}
for i := 0; i < len(added); i += 1 {
err = manager.newFromCloudSku(ctx, userCred, added[i])
err = region.newFromPublicCloudSku(ctx, userCred, added[i].GetExternalId())
if err != nil {
syncResult.AddError(err)
} else {
@@ -478,22 +434,58 @@ func (self *SElasticcacheSku) MarkAsSoldout(ctx context.Context) error {
return nil
})
return errors.Wrap(err, "ElasticcacheSku.MarkAsSoldout")
return errors.Wrap(err, "MarkAsSoldout")
}
func (self *SElasticcacheSku) syncWithCloudSku(ctx context.Context, userCred mcclient.TokenCredential, extSku SElasticcacheSku) error {
_, err := db.Update(self, func() error {
self.PrepaidStatus = extSku.PrepaidStatus
self.PostpaidStatus = extSku.PostpaidStatus
self.ZoneId = extSku.ZoneId
self.SlaveZoneId = extSku.SlaveZoneId
return nil
})
return err
}
func (manager *SElasticcacheSkuManager) newFromCloudSku(ctx context.Context, userCred mcclient.TokenCredential, extSku SElasticcacheSku) error {
return manager.TableSpec().Insert(ctx, &extSku)
func (self *SCloudregion) newFromPublicCloudSku(ctx context.Context, userCred mcclient.TokenCredential, externalId string) error {
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return err
}
zones, err := self.GetZones()
if err != nil {
return errors.Wrap(err, "GetZones")
}
zoneMaps := map[string]string{}
for _, zone := range zones {
zoneMaps[zone.ExternalId] = zone.Id
}
skuUrl := fmt.Sprintf("%s/%s/%s.json", meta.ElasticCacheBase, self.ExternalId, externalId)
sku := &SElasticcacheSku{}
sku.SetModelManager(ElasticcacheSkuManager, sku)
err = meta.Get(skuUrl, sku)
if err != nil {
return errors.Wrapf(err, "Get")
}
sku.Status = api.SkuStatusAvailable
sku.CloudregionId = self.Id
sku.Provider = self.Provider
if len(sku.ZoneId) > 0 {
zoneId := yunionmeta.GetZoneIdBySuffix(zoneMaps, sku.ZoneId)
if len(zoneId) == 0 {
return errors.Wrapf(err, "empty zoneId for %s", sku.ZoneId)
}
sku.ZoneId = zoneId
}
if len(sku.SlaveZoneId) > 0 {
zoneId := yunionmeta.GetZoneIdBySuffix(zoneMaps, sku.SlaveZoneId)
if len(zoneId) == 0 {
return errors.Wrapf(err, "empty zoneId for %s", sku.SlaveZoneId)
}
sku.SlaveZoneId = zoneId
}
return ElasticcacheSkuManager.TableSpec().Insert(ctx, sku)
}
func (manager *SElasticcacheSkuManager) GetPropertyInstanceSpecs(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) {
@@ -778,3 +770,77 @@ func (self *SCloudregion) newFromCloudElasticcacheSku(ctx context.Context, userC
return ElasticcacheSkuManager.TableSpec().Insert(ctx, sku)
}
// 全量同步elasticcache sku列表.
func SyncElasticCacheSkus(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
if isStart {
cnt, err := CloudaccountManager.Query().IsTrue("is_public_cloud").CountWithError()
if err != nil && err != sql.ErrNoRows {
log.Debugf("SyncElasticCacheSkus %s.sync skipped...", err)
return
} else if cnt == 0 {
log.Debugf("SyncElasticCacheSkus no public cloud.sync skipped...")
return
}
cnt, err = ElasticcacheSkuManager.Query().Limit(1).CountWithError()
if err != nil && err != sql.ErrNoRows {
log.Errorf("SyncElasticCacheSkus.QueryElasticcacheSku %s", err)
return
} else if cnt > 0 {
log.Debugf("SyncElasticCacheSkus synced skus, skip...")
return
}
}
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
log.Errorf("FetchYunionmeta %v", err)
return
}
index, err := meta.Index(ElasticcacheSkuManager.Keyword())
if err != nil {
log.Errorf("get cache sku index error: %v", err)
return
}
cloudregions := fetchSkuSyncCloudregions()
for i := range cloudregions {
region := &cloudregions[i]
if !region.GetDriver().IsSupportedElasticcache() {
continue
}
skuMeta := &SElasticcacheSku{}
skuMeta.SetModelManager(ElasticcacheSkuManager, skuMeta)
skuMeta.Id = region.ExternalId
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if !ok || newMd5 == yunionmeta.EMPTY_MD5 || len(oldMd5) > 0 && newMd5 == oldMd5 {
continue
}
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
result := ElasticcacheSkuManager.SyncElasticcacheSkus(ctx, userCred, region, false)
notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s result: %s", region.Name, result.Result())
log.Debugf(notes)
}
}
// 同步Region elasticcache sku列表.
func SyncElasticCacheSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, xor bool) error {
if !region.GetDriver().IsSupportedElasticcache() {
notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s not support elasticcache", region.Name)
log.Infof(notes)
return nil
}
result := ElasticcacheSkuManager.SyncElasticcacheSkus(ctx, userCred, region, xor)
notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s result: %s", region.Name, result.Result())
log.Infof(notes)
return nil
}
+2 -2
View File
@@ -490,8 +490,8 @@ func (fileSystem *SFileSystem) SyncWithCloudFileSystem(ctx context.Context, user
return nil
}
func (fileSystem *SCloudregion) getZoneIdBySuffix(zoneId string) (string, error) {
zones, err := fileSystem.GetZones()
func (region *SCloudregion) getZoneIdBySuffix(zoneId string) (string, error) {
zones, err := region.GetZones()
if err != nil {
return "", errors.Wrapf(err, "region.GetZones")
}
+66 -36
View File
@@ -32,6 +32,7 @@ import (
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/util/stringutils2"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
type SNasSkuManager struct {
@@ -198,22 +199,28 @@ func (self SNasSku) GetGlobalId() string {
return self.ExternalId
}
func (self *SCloudregion) SyncNasSkus(ctx context.Context, userCred mcclient.TokenCredential, meta *SSkuResourcesMeta, xor bool) compare.SyncResult {
lockman.LockRawObject(ctx, self.Id, "nas-sku")
defer lockman.ReleaseRawObject(ctx, self.Id, "nas-sku")
func (self *SCloudregion) SyncNasSkus(ctx context.Context, userCred mcclient.TokenCredential, xor bool) compare.SyncResult {
lockman.LockRawObject(ctx, self.Id, NasSkuManager.Keyword())
defer lockman.ReleaseRawObject(ctx, self.Id, NasSkuManager.Keyword())
syncResult := compare.SyncResult{}
result := compare.SyncResult{}
iskus, err := meta.GetNasSkusByRegionExternalId(self.ExternalId)
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(errors.Wrapf(err, "FetchYunionmeta"))
return result
}
iskus := []SNasSku{}
err = meta.List(NasSkuManager.Keyword(), self.ExternalId, &iskus)
if err != nil {
result.Error(err)
return result
}
dbSkus, err := self.GetNasSkus()
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(err)
return result
}
removed := make([]SNasSku, 0)
@@ -223,55 +230,86 @@ func (self *SCloudregion) SyncNasSkus(ctx context.Context, userCred mcclient.Tok
err = compare.CompareSets(dbSkus, iskus, &removed, &commondb, &commonext, &added)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(err)
return result
}
for i := 0; i < len(removed); i += 1 {
err = removed[i].Delete(ctx, userCred)
if err != nil {
syncResult.DeleteError(err)
result.DeleteError(err)
continue
}
syncResult.Delete()
result.Delete()
}
if !xor {
for i := 0; i < len(commondb); i += 1 {
err = commondb[i].syncWithCloudSku(ctx, userCred, commonext[i])
if err != nil {
syncResult.UpdateError(err)
result.UpdateError(err)
continue
}
syncResult.Update()
result.Update()
}
}
for i := 0; i < len(added); i += 1 {
err = self.newFromCloudNasSku(ctx, userCred, added[i])
if err != nil {
syncResult.AddError(err)
result.AddError(err)
} else {
syncResult.Add()
result.Add()
}
}
return syncResult
return result
}
func (self *SNasSku) syncWithCloudSku(ctx context.Context, userCred mcclient.TokenCredential, sku SNasSku) error {
_, err := db.Update(self, func() error {
jsonutils.Update(self, sku)
self.Status = api.NAS_SKU_AVAILABLE
self.PrepaidStatus = sku.PrepaidStatus
self.PostpaidStatus = sku.PostpaidStatus
return nil
})
return err
}
func (self *SCloudregion) newFromCloudNasSku(ctx context.Context, userCred mcclient.TokenCredential, isku SNasSku) error {
sku := &isku
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return err
}
zones, err := self.GetZones()
if err != nil {
return errors.Wrap(err, "GetZones")
}
zoneMaps := map[string]string{}
for _, zone := range zones {
zoneMaps[zone.ExternalId] = zone.Id
}
sku := &SNasSku{}
sku.SetModelManager(NasSkuManager, sku)
sku.Id = "" //避免使用yunion meta的id,导致出现duplicate entry问题
skuUrl := fmt.Sprintf("%s/%s/%s.json", meta.NasBase, self.ExternalId, isku.GetGlobalId())
err = meta.Get(skuUrl, sku)
if err != nil {
return errors.Wrapf(err, "Get")
}
if len(sku.ZoneIds) > 0 {
zoneIds := []string{}
for _, zoneExtId := range strings.Split(sku.ZoneIds, ",") {
zoneId := yunionmeta.GetZoneIdBySuffix(zoneMaps, zoneExtId) // Huawei rds sku zone1 maybe is cn-north-4f
if len(zoneId) > 0 {
zoneIds = append(zoneIds, zoneId)
}
}
sku.ZoneIds = strings.Join(zoneIds, ",")
}
sku.Status = api.NAS_SKU_AVAILABLE
sku.SetEnabled(true)
sku.CloudregionId = self.Id
sku.Provider = self.Provider
return NasSkuManager.TableSpec().Insert(ctx, sku)
}
@@ -309,12 +347,12 @@ func SyncRegionNasSkus(ctx context.Context, userCred mcclient.TokenCredential, r
return errors.Wrapf(err, "db.FetchModelObjects")
}
meta, err := FetchSkuResourcesMeta()
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return errors.Wrapf(err, "FetchSkuResourcesMeta")
return errors.Wrapf(err, "FetchYunionmeta")
}
index, err := meta.getSkuIndex(meta.NasBase)
index, err := meta.Index(NasSkuManager.Keyword())
if err != nil {
log.Errorf("get nas sku index error: %v", err)
return err
@@ -333,21 +371,13 @@ func SyncRegionNasSkus(ctx context.Context, userCred mcclient.TokenCredential, r
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
log.Debugf("%s nas skus is empty skip syncing", region.Name)
if !ok || newMd5 == yunionmeta.EMPTY_MD5 || len(oldMd5) > 0 && newMd5 == oldMd5 {
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Debugf("%s nas skus not changed skip syncing", region.Name)
continue
}
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
result := regions[i].SyncNasSkus(ctx, userCred, meta, xor)
result := regions[i].SyncNasSkus(ctx, userCred, xor)
msg := result.Result()
notes := fmt.Sprintf("SyncNasSkus for region %s result: %s", regions[i].Name, msg)
log.Debugf(notes)
+66 -35
View File
@@ -34,6 +34,7 @@ import (
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/util/stringutils2"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
type SNatSkuManager struct {
@@ -201,22 +202,29 @@ func (self SNatSku) GetGlobalId() string {
return self.ExternalId
}
func (self *SCloudregion) SyncNatSkus(ctx context.Context, userCred mcclient.TokenCredential, meta *SSkuResourcesMeta, xor bool) compare.SyncResult {
func (self *SCloudregion) SyncNatSkus(ctx context.Context, userCred mcclient.TokenCredential, xor bool) compare.SyncResult {
lockman.LockRawObject(ctx, self.Id, NatSkuManager.Keyword())
defer lockman.ReleaseRawObject(ctx, self.Id, NatSkuManager.Keyword())
syncResult := compare.SyncResult{}
result := compare.SyncResult{}
iskus, err := meta.GetNatSkusByRegionExternalId(self.ExternalId)
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(errors.Wrapf(err, "FetchYunionmeta"))
return result
}
iskus := []SNatSku{}
err = meta.List(NatSkuManager.Keyword(), self.ExternalId, &iskus)
if err != nil {
result.Error(err)
return result
}
dbSkus, err := self.GetNatSkus()
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(err)
return result
}
removed := make([]SNatSku, 0)
@@ -226,55 +234,86 @@ func (self *SCloudregion) SyncNatSkus(ctx context.Context, userCred mcclient.Tok
err = compare.CompareSets(dbSkus, iskus, &removed, &commondb, &commonext, &added)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(err)
return result
}
for i := 0; i < len(removed); i += 1 {
err = removed[i].Delete(ctx, userCred)
if err != nil {
syncResult.DeleteError(err)
result.DeleteError(err)
continue
}
syncResult.Delete()
result.Delete()
}
if !xor {
for i := 0; i < len(commondb); i += 1 {
err = commondb[i].syncWithCloudSku(ctx, userCred, commonext[i])
if err != nil {
syncResult.UpdateError(err)
result.UpdateError(err)
continue
}
syncResult.Update()
result.Update()
}
}
for i := 0; i < len(added); i += 1 {
err = self.newFromCloudNatSku(ctx, userCred, added[i])
if err != nil {
syncResult.AddError(err)
result.AddError(err)
} else {
syncResult.Add()
result.Add()
}
}
return syncResult
return result
}
func (self *SNatSku) syncWithCloudSku(ctx context.Context, userCred mcclient.TokenCredential, sku SNatSku) error {
_, err := db.Update(self, func() error {
jsonutils.Update(self, sku)
self.Status = api.NAT_SKU_AVAILABLE
self.PrepaidStatus = sku.PrepaidStatus
self.PostpaidStatus = sku.PostpaidStatus
return nil
})
return err
}
func (self *SCloudregion) newFromCloudNatSku(ctx context.Context, userCred mcclient.TokenCredential, isku SNatSku) error {
sku := &isku
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return err
}
zones, err := self.GetZones()
if err != nil {
return errors.Wrap(err, "GetZones")
}
zoneMaps := map[string]string{}
for _, zone := range zones {
zoneMaps[zone.ExternalId] = zone.Id
}
sku := &SNatSku{}
sku.SetModelManager(NatSkuManager, sku)
sku.Id = "" //避免使用yunion meta的id,导致出现duplicate entry问题
sku.Status = api.NAT_SKU_AVAILABLE
skuUrl := fmt.Sprintf("%s/%s/%s.json", meta.NatBase, self.ExternalId, isku.GetGlobalId())
err = meta.Get(skuUrl, sku)
if err != nil {
return errors.Wrapf(err, "Get")
}
if len(sku.ZoneIds) > 0 {
zoneIds := []string{}
for _, zoneExtId := range strings.Split(sku.ZoneIds, ",") {
zoneId := yunionmeta.GetZoneIdBySuffix(zoneMaps, zoneExtId) // Huawei rds sku zone1 maybe is cn-north-4f
if len(zoneId) > 0 {
zoneIds = append(zoneIds, zoneId)
}
}
sku.ZoneIds = strings.Join(zoneIds, ",")
}
sku.Status = api.NAS_SKU_AVAILABLE
sku.SetEnabled(true)
sku.CloudregionId = self.Id
sku.Provider = self.Provider
return NatSkuManager.TableSpec().Insert(ctx, sku)
}
@@ -312,12 +351,12 @@ func SyncRegionNatSkus(ctx context.Context, userCred mcclient.TokenCredential, r
return errors.Wrapf(err, "db.FetchModelObjects")
}
meta, err := FetchSkuResourcesMeta()
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return errors.Wrapf(err, "FetchSkuResourcesMeta")
return errors.Wrapf(err, "FetchYunionmeta")
}
index, err := meta.getSkuIndex(meta.NatBase)
index, err := meta.Index(NatSkuManager.Keyword())
if err != nil {
log.Errorf("get nat sku index error: %v", err)
return err
@@ -336,21 +375,13 @@ func SyncRegionNatSkus(ctx context.Context, userCred mcclient.TokenCredential, r
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
log.Infof("%s Nat Skus is empty skip syncing", region.Name)
if !ok || newMd5 == yunionmeta.EMPTY_MD5 || len(oldMd5) > 0 && newMd5 == oldMd5 {
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Infof("%s Nat Skus not Changed skip syncing", region.Name)
continue
}
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
result := regions[i].SyncNatSkus(ctx, userCred, meta, xor)
result := regions[i].SyncNatSkus(ctx, userCred, xor)
msg := result.Result()
notes := fmt.Sprintf("SyncNatSkus for region %s result: %s", regions[i].Name, msg)
log.Infof(notes)
@@ -44,6 +44,7 @@ import (
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/mcclient/modules/scheduler"
"yunion.io/x/onecloud/pkg/util/stringutils2"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
type SServerSkuManager struct {
@@ -1074,7 +1075,7 @@ func (manager *SServerSkuManager) SyncPrivateCloudSkus(
}
for i := 0; i < len(added); i += 1 {
err := manager.newFromCloudSku(ctx, userCred, region, added[i])
err := manager.newPrivateCloudSku(ctx, userCred, region, added[i])
if err != nil {
result.AddError(err)
continue
@@ -1145,7 +1146,45 @@ func (self *SServerSku) setPrepaidPostpaidStatus(userCred mcclient.TokenCredenti
return nil
}
func (manager *SServerSkuManager) newFromCloudSku(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSku cloudprovider.ICloudSku) error {
func (region *SCloudregion) newPublicCloudSku(ctx context.Context, userCred mcclient.TokenCredential, extSku SServerSku) error {
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return err
}
zones, err := region.GetZones()
if err != nil {
return errors.Wrap(err, "GetZones")
}
zoneMaps := map[string]string{}
for _, zone := range zones {
zoneMaps[zone.ExternalId] = zone.Id
}
sku := &SServerSku{}
sku.SetModelManager(ServerSkuManager, sku)
skuUrl := fmt.Sprintf("%s/%s/%s.json", meta.ServerBase, region.ExternalId, extSku.ExternalId)
err = meta.Get(skuUrl, sku)
if err != nil {
return errors.Wrapf(err, "Get")
}
if len(sku.ZoneId) > 0 {
zoneId := yunionmeta.GetZoneIdBySuffix(zoneMaps, sku.ZoneId)
if len(zoneId) > 0 {
sku.ZoneId = zoneId
}
}
// 第一次同步新建的套餐是启用状态
sku.Enabled = tristate.True
sku.Status = api.SkuStatusReady
sku.CloudregionId = region.Id
sku.Provider = region.Provider
return ServerSkuManager.TableSpec().Insert(ctx, sku)
}
func (manager *SServerSkuManager) newPrivateCloudSku(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSku cloudprovider.ICloudSku) error {
sku := &SServerSku{Provider: region.Provider}
sku.SetModelManager(manager, sku)
@@ -1157,33 +1196,13 @@ func (manager *SServerSkuManager) newFromCloudSku(ctx context.Context, userCred
sku.CloudregionId = region.Id
sku.Name = extSku.GetName()
sku.SetModelManager(manager, sku)
err := manager.TableSpec().Insert(ctx, sku)
if err != nil {
return errors.Wrapf(err, "Insert")
}
db.OpsLog.LogEvent(sku, db.ACT_CREATE, sku.GetShortDesc(ctx), userCred)
return nil
}
func (manager *SServerSkuManager) newPublicCloudSku(ctx context.Context, userCred mcclient.TokenCredential, extSku SServerSku) error {
extSku.Enabled = tristate.True
extSku.Status = api.SkuStatusReady
return manager.TableSpec().Insert(ctx, &extSku)
return manager.TableSpec().Insert(ctx, sku)
}
func (self *SServerSku) syncWithCloudSku(ctx context.Context, userCred mcclient.TokenCredential, extSku SServerSku) error {
_, err := db.Update(self, func() error {
self.ZoneId = extSku.ZoneId
self.InstanceTypeCategory = extSku.InstanceTypeCategory
self.PrepaidStatus = extSku.PrepaidStatus
self.PostpaidStatus = extSku.PostpaidStatus
self.Name = extSku.GetName()
self.CpuArch = extSku.CpuArch
self.SysDiskType = extSku.SysDiskType
self.DataDiskTypes = extSku.DataDiskTypes
return nil
})
return err
@@ -1212,22 +1231,29 @@ func (manager *SServerSkuManager) FetchSkusByRegion(regionID string) ([]SServerS
return skus, nil
}
func (manager *SServerSkuManager) SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSkuMeta *SSkuResourcesMeta, xor bool) compare.SyncResult {
func (manager *SServerSkuManager) SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, xor bool) compare.SyncResult {
lockman.LockRawObject(ctx, manager.Keyword(), region.Id)
defer lockman.ReleaseRawObject(ctx, manager.Keyword(), region.Id)
syncResult := compare.SyncResult{}
result := compare.SyncResult{}
extSkus, err := extSkuMeta.GetServerSkusByRegionExternalId(region.ExternalId)
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(errors.Wrapf(err, "FetchYunionmeta"))
return result
}
extSkus := []SServerSku{}
err = meta.List(manager.Keyword(), region.ExternalId, &extSkus)
if err != nil {
result.Error(errors.Wrapf(err, "List"))
return result
}
dbSkus, err := manager.FetchSkusByRegion(region.GetId())
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(err)
return result
}
removed := make([]SServerSku, 0)
@@ -1237,8 +1263,8 @@ func (manager *SServerSkuManager) SyncServerSkus(ctx context.Context, userCred m
err = compare.CompareSets(dbSkus, extSkus, &removed, &commondb, &commonext, &added)
if err != nil {
syncResult.Error(err)
return syncResult
result.Error(err)
return result
}
for i := 0; i < len(removed); i += 1 {
@@ -1249,27 +1275,27 @@ func (manager *SServerSkuManager) SyncServerSkus(ctx context.Context, userCred m
err = removed[i].RealDelete(ctx, userCred)
}
if err != nil {
syncResult.DeleteError(err)
result.DeleteError(err)
} else {
syncResult.Delete()
result.Delete()
}
}
if !xor {
for i := 0; i < len(commondb); i += 1 {
err = commondb[i].syncWithCloudSku(ctx, userCred, commonext[i])
if err != nil {
syncResult.UpdateError(err)
result.UpdateError(err)
} else {
syncResult.Update()
result.Update()
}
}
}
for i := 0; i < len(added); i += 1 {
err = manager.newPublicCloudSku(ctx, userCred, added[i])
err = region.newPublicCloudSku(ctx, userCred, added[i])
if err != nil {
syncResult.AddError(err)
result.AddError(err)
} else {
syncResult.Add()
result.Add()
}
}
@@ -1279,7 +1305,7 @@ func (manager *SServerSkuManager) SyncServerSkus(ctx context.Context, userCred m
log.Errorf("SchedManager SyncSku %s", err)
}
return syncResult
return result
}
// sku标记为soldout状态。
@@ -1503,3 +1529,77 @@ func (self *SServerSku) GetICloudSku(ctx context.Context) (cloudprovider.ICloudS
}
return nil, errors.Wrapf(cloudprovider.ErrNotFound, self.ExternalId)
}
func fetchSkuSyncCloudregions() []SCloudregion {
cloudregions := []SCloudregion{}
q := CloudregionManager.Query()
q = q.In("provider", CloudproviderManager.GetPublicProviderProvidersQuery())
err := db.FetchModelObjects(CloudregionManager, q, &cloudregions)
if err != nil {
log.Errorf("fetchSkuSyncCloudregions.FetchCloudregions failed: %v", err)
return nil
}
return cloudregions
}
// 全量同步sku列表.
func SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
if isStart {
cnt, err := ServerSkuManager.GetPublicCloudSkuCount()
if err != nil {
log.Errorf("GetPublicCloudSkuCount fail %s", err)
return
}
if cnt > 0 {
log.Debugf("GetPublicCloudSkuCount synced skus, skip...")
return
}
}
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
log.Errorf("FetchYunionmeta %v", err)
return
}
index, err := meta.Index(ServerSkuManager.Keyword())
if err != nil {
log.Errorf("getServerSkuIndex error: %v", err)
return
}
cloudregions := fetchSkuSyncCloudregions()
for i := range cloudregions {
region := &cloudregions[i]
skuMeta := &SServerSku{}
skuMeta.SetModelManager(ServerSkuManager, skuMeta)
skuMeta.Id = region.ExternalId
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if !ok || newMd5 == yunionmeta.EMPTY_MD5 || len(oldMd5) > 0 && newMd5 == oldMd5 {
continue
}
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
result := ServerSkuManager.SyncServerSkus(ctx, userCred, region, false)
notes := fmt.Sprintf("SyncServerSkusByRegion %s result: %s", region.Name, result.Result())
log.Debugf(notes)
}
// 清理无效的sku
log.Debugf("DeleteInvalidSkus in processing...")
ServerSkuManager.PendingDeleteInvalidSku()
}
// 同步指定region sku列表
func SyncServerSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, xor bool) compare.SyncResult {
result := compare.SyncResult{}
result = ServerSkuManager.SyncServerSkus(ctx, userCred, region, xor)
notes := fmt.Sprintf("SyncServerSkusByRegion %s result: %s", region.Name, result.Result())
log.Infof(notes)
return result
}
-695
View File
@@ -1,695 +0,0 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package models
import (
"context"
"database/sql"
"fmt"
"net/http"
"strings"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/compare"
"yunion.io/x/pkg/util/httputils"
v "yunion.io/x/pkg/util/version"
"yunion.io/x/pkg/utils"
apis "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/compute/options"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/mcclient/modules/compute"
)
const (
EMPTY_MD5 = "d751713988987e9331980363e24189ce"
)
/*
资源套餐下载连接信息
server: 虚拟机
elasticcache: 弹性缓存(redis&memcached)
*/
type SSkuResourcesMeta struct {
DBInstanceBase string `json:"dbinstance_base"`
ServerBase string `json:"server_base"`
ElasticCacheBase string `json:"elastic_cache_base"`
ImageBase string `json:"image_base"`
NatBase string `json:"nat_base"`
NasBase string `json:"nas_base"`
WafBase string `json:"waf_base"`
}
func (self *SSkuResourcesMeta) getZoneIdBySuffix(zoneMaps map[string]string, suffix string) string {
for externalId, id := range zoneMaps {
if strings.HasSuffix(externalId, suffix) {
return id
}
}
return ""
}
func (self *SSkuResourcesMeta) GetCloudimages(regionExternalId string) ([]SCachedimage, error) {
objs, err := self.getObjsByRegion(self.ImageBase, regionExternalId)
if err != nil {
return nil, errors.Wrapf(err, "getObjsByRegion")
}
images := []SCachedimage{}
err = jsonutils.Update(&images, objs)
if err != nil {
return nil, errors.Wrapf(err, "jsonutils.Update")
}
return images, nil
}
func (self *SSkuResourcesMeta) GetDBInstanceSkusByRegionExternalId(regionExternalId string) ([]SDBInstanceSku, error) {
regionId, zoneMaps, err := self.GetRegionIdAndZoneMaps(regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "GetRegionIdAndZoneMaps")
}
result := []SDBInstanceSku{}
objs, err := self.getObjsByRegion(self.DBInstanceBase, regionExternalId)
if err != nil {
return nil, errors.Wrapf(err, "getSkusByRegion")
}
noZoneIds, cnt := []string{}, 0
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) // Huawei rds sku zone1 maybe is cn-north-4f
if len(zoneId) == 0 {
if !utils.IsInStringArray(sku.Zone1, noZoneIds) {
noZoneIds = append(noZoneIds, sku.Zone1)
}
cnt++
continue
}
sku.Zone1 = zoneId
}
if len(sku.Zone2) > 0 {
zoneId := self.getZoneIdBySuffix(zoneMaps, sku.Zone2)
if len(zoneId) == 0 {
if !utils.IsInStringArray(sku.Zone2, noZoneIds) {
noZoneIds = append(noZoneIds, sku.Zone2)
}
cnt++
continue
}
sku.Zone2 = zoneId
}
if len(sku.Zone3) > 0 {
zoneId := self.getZoneIdBySuffix(zoneMaps, sku.Zone3)
if len(zoneId) == 0 {
if !utils.IsInStringArray(sku.Zone3, noZoneIds) {
noZoneIds = append(noZoneIds, sku.Zone3)
}
cnt++
continue
}
sku.Zone3 = zoneId
}
sku.Id = ""
sku.CloudregionId = regionId
result = append(result, sku)
}
if len(noZoneIds) > 0 {
log.Warningf("can not fetch rds sku %d zone %s for %s", cnt, noZoneIds, regionExternalId)
}
return result, nil
}
func (self *SSkuResourcesMeta) GetNatSkusByRegionExternalId(regionExternalId string) ([]SNatSku, error) {
regionId, zoneMaps, err := self.GetRegionIdAndZoneMaps(regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "GetRegionIdAndZoneMaps")
}
result := []SNatSku{}
objs, err := self.getObjsByRegion(self.NatBase, regionExternalId)
if err != nil {
return nil, errors.Wrapf(err, "getSkusByRegion")
}
noZoneIds, cnt := []string{}, 0
for _, obj := range objs {
sku := SNatSku{}
sku.SetModelManager(NatSkuManager, &sku)
err = obj.Unmarshal(&sku)
if err != nil {
return nil, errors.Wrapf(err, "obj.Unmarshal")
}
if len(sku.ZoneIds) > 0 {
zoneIds := []string{}
for _, zoneExtId := range strings.Split(sku.ZoneIds, ",") {
zoneId := self.getZoneIdBySuffix(zoneMaps, zoneExtId) // Huawei rds sku zone1 maybe is cn-north-4f
if len(zoneId) == 0 {
if !utils.IsInStringArray(zoneExtId, noZoneIds) {
noZoneIds = append(noZoneIds, zoneExtId)
}
cnt++
continue
}
zoneIds = append(zoneIds, zoneId)
}
sku.ZoneIds = strings.Join(zoneIds, ",")
}
sku.Id = ""
sku.CloudregionId = regionId
result = append(result, sku)
}
if len(noZoneIds) > 0 {
log.Warningf("can not fetch nat sku %d zone %s for %s", cnt, noZoneIds, regionExternalId)
}
return result, nil
}
func (self *SSkuResourcesMeta) GetNasSkusByRegionExternalId(regionExternalId string) ([]SNasSku, error) {
regionId, zoneMaps, err := self.GetRegionIdAndZoneMaps(regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "GetRegionIdAndZoneMaps")
}
result := []SNasSku{}
objs, err := self.getObjsByRegion(self.NasBase, regionExternalId)
if err != nil {
return nil, errors.Wrapf(err, "getSkusByRegion")
}
noZoneIds, cnt := []string{}, 0
for _, obj := range objs {
sku := SNasSku{}
sku.SetModelManager(NasSkuManager, &sku)
err = obj.Unmarshal(&sku)
if err != nil {
return nil, errors.Wrapf(err, "obj.Unmarshal")
}
if len(sku.ZoneIds) > 0 {
zoneIds := []string{}
for _, zoneExtId := range strings.Split(sku.ZoneIds, ",") {
zoneId := self.getZoneIdBySuffix(zoneMaps, zoneExtId) // Huawei rds sku zone1 maybe is cn-north-4f
if len(zoneId) == 0 {
if !utils.IsInStringArray(zoneExtId, noZoneIds) {
noZoneIds = append(noZoneIds, zoneExtId)
}
cnt++
continue
}
zoneIds = append(zoneIds, zoneId)
}
sku.ZoneIds = strings.Join(zoneIds, ",")
}
sku.Id = ""
sku.CloudregionId = regionId
result = append(result, sku)
}
if len(noZoneIds) > 0 {
log.Warningf("can not fetch nas sku %d zone %s for %s", cnt, noZoneIds, regionExternalId)
}
return result, nil
}
func (self *SSkuResourcesMeta) getCloudregion(regionExternalId string) (*SCloudregion, error) {
region, err := db.FetchByExternalId(CloudregionManager, regionExternalId)
if err != nil {
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.getObjsByRegion(self.ServerBase, regionExternalId)
if err != nil {
return nil, errors.Wrap(err, "getSkusByRegion")
}
noZoneIds, cnt := []string{}, 0
for _, obj := range objs {
sku := SServerSku{}
sku.SetModelManager(ServerSkuManager, &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 {
if !utils.IsInStringArray(sku.ZoneId, noZoneIds) {
noZoneIds = append(noZoneIds, sku.ZoneId)
}
cnt++
continue
}
sku.ZoneId = zoneId
}
sku.Id = ""
sku.CloudregionId = regionId
result = append(result, sku)
}
if len(noZoneIds) > 0 {
log.Warningf("can not fetch server sku %d zone id %s for region %s", cnt, noZoneIds, regionExternalId)
}
return result, nil
}
func getElaticCacheSkuRegionExtId(regionExtId string) string {
if strings.HasPrefix(regionExtId, apis.CLOUD_ACCESS_ENV_ALIYUN_FINANCE) && strings.HasSuffix(regionExtId, "cn-hangzhou") {
return regionExtId + "-finance"
}
return regionExtId
}
func getElaticCacheSkuZoneId(zoneExtId string) string {
if strings.HasPrefix(zoneExtId, apis.CLOUD_ACCESS_ENV_ALIYUN_FINANCE) {
if strings.HasSuffix(zoneExtId, "cn-hangzhou-finance-b") || strings.HasSuffix(zoneExtId, "cn-hangzhou-finance-c") || strings.HasSuffix(zoneExtId, "cn-hangzhou-finance-d") {
zoneExtId = strings.Replace(zoneExtId, "-finance", "", -1)
} else if strings.Contains(zoneExtId, "cn-hangzhou") {
zoneExtId = strings.Replace(zoneExtId, "-finance", "", 1)
}
}
return zoneExtId
}
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{}
noZoneIds, cnt := []string{}, 0
// aliyun finance cloud
remoteRegion := getElaticCacheSkuRegionExtId(regionExternalId)
objs, err := self.getObjsByRegion(self.ElasticCacheBase, remoteRegion)
if err != nil {
return nil, errors.Wrap(err, "getObjsByRegion")
}
for _, obj := range objs {
sku := SElasticcacheSku{}
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, getElaticCacheSkuZoneId(sku.ZoneId))
if len(zoneId) == 0 {
if !utils.IsInStringArray(sku.ZoneId, noZoneIds) {
noZoneIds = append(noZoneIds, sku.ZoneId)
}
cnt++
continue
}
sku.ZoneId = zoneId
}
if len(sku.SlaveZoneId) > 0 {
zoneId := self.getZoneIdBySuffix(zoneMaps, getElaticCacheSkuZoneId(sku.SlaveZoneId))
if len(zoneId) == 0 {
if !utils.IsInStringArray(sku.SlaveZoneId, noZoneIds) {
noZoneIds = append(noZoneIds, sku.SlaveZoneId)
}
cnt++
continue
}
sku.SlaveZoneId = zoneId
}
sku.Id = ""
sku.CloudregionId = regionId
result = append(result, sku)
}
if len(noZoneIds) > 0 {
log.Warningf("can not fetch redis sku %d zone %s for %s", cnt, noZoneIds, regionExternalId)
}
return result, nil
}
func (self *SSkuResourcesMeta) getObjsByRegion(base string, region string) ([]jsonutils.JSONObject, error) {
url := fmt.Sprintf("%s/%s.json", base, region)
items, err := self._get(url)
if err != nil {
return nil, errors.Wrap(err, "getSkusByRegion.get")
}
return items, nil
}
func (self *SSkuResourcesMeta) request(url string) (jsonutils.JSONObject, error) {
client := httputils.GetAdaptiveTimeoutClient()
header := http.Header{}
header.Set("User-Agent", "vendor/yunion-OneCloud@"+v.Get().GitVersion)
_, resp, err := httputils.JSONRequest(client, context.TODO(), httputils.GET, url, header, nil, false)
return resp, err
}
func (self *SSkuResourcesMeta) getServerSkuIndex() (map[string]string, error) {
resp, err := self.request(fmt.Sprintf("%s/index.json", self.ServerBase))
if err != nil {
return map[string]string{}, errors.Wrapf(err, "request")
}
ret := map[string]string{}
err = resp.Unmarshal(ret)
if err != nil {
return map[string]string{}, errors.Wrapf(err, "resp.Unmarshal")
}
return ret, nil
}
func (self *SSkuResourcesMeta) getSkuIndex(res string) (map[string]string, error) {
resp, err := self.request(fmt.Sprintf("%s/index.json", res))
if err != nil {
return map[string]string{}, errors.Wrapf(err, "request")
}
ret := map[string]string{}
err = resp.Unmarshal(ret)
if err != nil {
return map[string]string{}, errors.Wrapf(err, "resp.Unmarshal")
}
return ret, nil
}
func (self *SSkuResourcesMeta) getCloudimageIndex() (map[string]string, error) {
resp, err := self.request(fmt.Sprintf("%s/index.json", self.ImageBase))
if err != nil {
return map[string]string{}, errors.Wrapf(err, "request")
}
ret := map[string]string{}
err = resp.Unmarshal(ret)
if err != nil {
return map[string]string{}, errors.Wrapf(err, "resp.Unmarshal")
}
return ret, nil
}
func (self *SSkuResourcesMeta) getWafIndex() (map[string]string, error) {
resp, err := self.request(fmt.Sprintf("%s/index.json", self.WafBase))
if err != nil {
return map[string]string{}, errors.Wrapf(err, "request")
}
ret := map[string]string{}
err = resp.Unmarshal(ret)
if err != nil {
return map[string]string{}, errors.Wrapf(err, "resp.Unmarshal")
}
return ret, nil
}
func (self *SSkuResourcesMeta) _get(url string) ([]jsonutils.JSONObject, error) {
if !strings.HasPrefix(url, "http") {
return nil, fmt.Errorf("SkuResourcesMeta.get invalid url %s.expected has prefix 'http'", url)
}
jsonContent, err := self.request(url)
if err != nil {
return nil, errors.Wrapf(err, "request %s", url)
}
var ret []jsonutils.JSONObject
err = jsonContent.Unmarshal(&ret)
if err != nil {
return nil, fmt.Errorf("SkuResourcesMeta.get.Unmarshal %s content: %s url: %s", err, jsonContent, url)
}
return ret, nil
}
// 全量同步elasticcache sku列表.
func SyncElasticCacheSkus(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
if isStart {
cnt, err := CloudaccountManager.Query().IsTrue("is_public_cloud").CountWithError()
if err != nil && err != sql.ErrNoRows {
log.Debugf("SyncElasticCacheSkus %s.sync skipped...", err)
return
} else if cnt == 0 {
log.Debugf("SyncElasticCacheSkus no public cloud.sync skipped...")
return
}
cnt, err = ElasticcacheSkuManager.Query().Limit(1).CountWithError()
if err != nil && err != sql.ErrNoRows {
log.Errorf("SyncElasticCacheSkus.QueryElasticcacheSku %s", err)
return
} else if cnt > 0 {
log.Debugf("SyncElasticCacheSkus synced skus, skip...")
return
}
}
meta, err := FetchSkuResourcesMeta()
if err != nil {
log.Errorf("SyncElasticCacheSkus.FetchSkuResourcesMeta %s", err)
return
}
index, err := meta.getSkuIndex(meta.ElasticCacheBase)
if err != nil {
log.Errorf("get cache sku index error: %v", err)
return
}
cloudregions := fetchSkuSyncCloudregions()
for i := range cloudregions {
region := &cloudregions[i]
if !region.GetDriver().IsSupportedElasticcache() {
continue
}
skuMeta := &SElasticcacheSku{}
skuMeta.SetModelManager(ElasticcacheSkuManager, skuMeta)
skuMeta.Id = region.ExternalId
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
log.Debugf("%s redis skus is empty skip syncing", region.Name)
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Debugf("%s redis skus not changed skip syncing", region.Name)
continue
}
result := ElasticcacheSkuManager.SyncElasticcacheSkus(ctx, userCred, region, meta, false)
notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s result: %s", region.Name, result.Result())
log.Debugf(notes)
}
}
// 同步Region elasticcache sku列表.
func SyncElasticCacheSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, xor bool) 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 {
return errors.Wrap(err, "SyncElasticCacheSkusByRegion.FetchSkuResourcesMeta")
}
result := ElasticcacheSkuManager.SyncElasticcacheSkus(ctx, userCred, region, meta, xor)
notes := fmt.Sprintf("SyncElasticCacheSkusByRegion %s result: %s", region.Name, result.Result())
log.Infof(notes)
return nil
}
// 全量同步sku列表.
func SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
if isStart {
cnt, err := ServerSkuManager.GetPublicCloudSkuCount()
if err != nil {
log.Errorf("GetPublicCloudSkuCount fail %s", err)
return
}
if cnt > 0 {
log.Debugf("GetPublicCloudSkuCount synced skus, skip...")
return
}
}
meta, err := FetchSkuResourcesMeta()
if err != nil {
log.Errorf("SyncServerSkus.FetchSkuResourcesMeta %s", err)
return
}
index, err := meta.getServerSkuIndex()
if err != nil {
log.Errorf("getServerSkuIndex error: %v", err)
return
}
cloudregions := fetchSkuSyncCloudregions()
for i := range cloudregions {
region := &cloudregions[i]
skuMeta := &SServerSku{}
skuMeta.SetModelManager(ServerSkuManager, skuMeta)
skuMeta.Id = region.ExternalId
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
log.Debugf("%s server skus is empty skip syncing", region.Name)
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Debugf("%s server skus not changed skip syncing", region.Name)
continue
}
result := ServerSkuManager.SyncServerSkus(ctx, userCred, region, meta, false)
notes := fmt.Sprintf("SyncServerSkusByRegion %s result: %s", region.Name, result.Result())
log.Debugf(notes)
}
// 清理无效的sku
log.Debugf("DeleteInvalidSkus in processing...")
ServerSkuManager.PendingDeleteInvalidSku()
}
// 同步指定region sku列表
func SyncServerSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSkuMeta *SSkuResourcesMeta, xor bool) compare.SyncResult {
result := compare.SyncResult{}
var err error
if extSkuMeta == nil {
extSkuMeta, err = FetchSkuResourcesMeta()
if err != nil {
result.AddError(errors.Wrap(err, "SyncServerSkusByRegion.FetchSkuResourcesMeta"))
return result
}
}
result = ServerSkuManager.SyncServerSkus(ctx, userCred, region, extSkuMeta, xor)
notes := fmt.Sprintf("SyncServerSkusByRegion %s result: %s", region.Name, result.Result())
log.Infof(notes)
return result
}
func FetchSkuResourcesMeta() (*SSkuResourcesMeta, error) {
s := auth.GetAdminSession(context.Background(), options.Options.Region)
transport := httputils.GetTransport(true)
transport.Proxy = options.Options.HttpTransportProxyFunc()
client := &http.Client{Transport: transport}
meta, err := compute.OfflineCloudmeta.GetSkuSourcesMeta(s, client)
if err != nil {
return nil, errors.Wrap(err, "fetchSkuSourceUrls.GetSkuSourcesMeta")
}
ret := &SSkuResourcesMeta{}
err = meta.Unmarshal(ret)
if err != nil {
return nil, errors.Wrap(err, "fetchSkuSourceUrls.Unmarshal")
}
return ret, nil
}
func fetchCloudEnvs() ([]string, error) {
accounts := []SCloudaccount{}
q := CloudaccountManager.Query("provider", "access_url").In("provider", CloudproviderManager.GetPublicProviderProvidersQuery()).Distinct()
err := q.All(&accounts)
if err != nil {
return nil, errors.Wrapf(err, "q.All")
}
ret := []string{}
for i := range accounts {
ret = append(ret, apis.GetCloudEnv(accounts[i].Provider, accounts[i].AccessUrl))
}
return ret, nil
}
func fetchSkuSyncCloudregions() []SCloudregion {
cloudregions := []SCloudregion{}
q := CloudregionManager.Query()
q = q.In("provider", CloudproviderManager.GetPublicProviderProvidersQuery())
err := db.FetchModelObjects(CloudregionManager, q, &cloudregions)
if err != nil {
log.Errorf("fetchSkuSyncCloudregions.FetchCloudregions failed: %v", err)
return nil
}
return cloudregions
}
type sWafGroup struct {
SWafRuleGroup
Rules []SWafRule
}
func (self sWafGroup) GetGlobalId() string {
return self.ExternalId
}
func (self SWafRule) GetGlobalId() string {
return self.ExternalId
}
func (self *SSkuResourcesMeta) getCloudWafGroups(cloudEnv string) ([]sWafGroup, error) {
url := fmt.Sprintf("%s/%s.json", self.WafBase, cloudEnv)
resp, err := self.request(url)
if err != nil {
return nil, errors.Wrapf(err, "_get(%s)", url)
}
ret := []sWafGroup{}
err = resp.Unmarshal(&ret)
if err != nil {
return nil, errors.Wrapf(err, "resp.Unmarshal")
}
return ret, nil
}
+55 -25
View File
@@ -30,6 +30,7 @@ import (
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/util/stringutils2"
"yunion.io/x/onecloud/pkg/util/yunionmeta"
)
type SWafRuleGroupManager struct {
@@ -161,8 +162,8 @@ func (self *SWafRuleGroup) RealDelete(ctx context.Context, userCred mcclient.Tok
return self.SStatusInfrasResourceBase.Delete(ctx, userCred)
}
func (self *SSkuResourcesMeta) GetWafGroups(cloudEnv string) ([]SWafRuleGroup, error) {
q := WafRuleGroupManager.Query().Equals("cloud_env", cloudEnv).IsTrue("is_system")
func (manager *SWafRuleGroupManager) GetWafGroups(cloudEnv string) ([]SWafRuleGroup, error) {
q := manager.Query().Equals("cloud_env", cloudEnv).IsTrue("is_system")
groups := []SWafRuleGroup{}
err := db.FetchModelObjects(WafRuleGroupManager, q, &groups)
return groups, err
@@ -187,9 +188,9 @@ func (self *SWafRuleGroup) syncWithCloudSku(ctx context.Context, userCred mcclie
return nil
}
func (self *SSkuResourcesMeta) newFromCloudWafGroup(ctx context.Context, userCred mcclient.TokenCredential, ext sWafGroup) error {
func (manager *SWafRuleGroupManager) newFromCloudWafGroup(ctx context.Context, userCred mcclient.TokenCredential, ext sWafGroup) error {
group := &ext.SWafRuleGroup
group.SetModelManager(WafRuleGroupManager, group)
group.SetModelManager(manager, group)
group.Status = api.WAF_RULE_GROUP_STATUS_AVAILABLE
group.IsPublic = true
err := WafRuleGroupManager.TableSpec().Insert(ctx, group)
@@ -204,17 +205,38 @@ func (self *SSkuResourcesMeta) newFromCloudWafGroup(ctx context.Context, userCre
return nil
}
func (self *SSkuResourcesMeta) SyncWafGroups(ctx context.Context, userCred mcclient.TokenCredential, cloudEnv string, isStart bool) compare.SyncResult {
lockman.LockRawObject(ctx, cloudEnv, "waf-rule-group")
defer lockman.ReleaseRawObject(ctx, cloudEnv, "waf-rule-group")
type sWafGroup struct {
SWafRuleGroup
Rules []SWafRule
}
func (self sWafGroup) GetGlobalId() string {
return self.ExternalId
}
func (self SWafRule) GetGlobalId() string {
return self.ExternalId
}
func (manager *SWafRuleGroupManager) SyncWafGroups(ctx context.Context, userCred mcclient.TokenCredential, cloudEnv string, isStart bool) compare.SyncResult {
lockman.LockRawObject(ctx, cloudEnv, manager.Keyword())
defer lockman.ReleaseRawObject(ctx, cloudEnv, manager.Keyword())
result := compare.SyncResult{}
exts, err := self.getCloudWafGroups(cloudEnv)
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
result.Error(errors.Wrapf(err, "getWafGroups(%s)", cloudEnv))
result.Error(errors.Wrapf(err, "FetchYunionmeta"))
return result
}
dbGroup, err := self.GetWafGroups(cloudEnv)
exts := []sWafGroup{}
err = meta.List(WafRuleManager.Keyword(), cloudEnv, &exts)
if err != nil {
result.Error(errors.Wrapf(err, "List(%s)", cloudEnv))
return result
}
dbGroup, err := manager.GetWafGroups(cloudEnv)
if err != nil {
result.Error(errors.Wrapf(err, "GetWafGroups"))
return result
@@ -253,7 +275,7 @@ func (self *SSkuResourcesMeta) SyncWafGroups(ctx context.Context, userCred mccli
result.Update()
}
for i := 0; i < len(added); i += 1 {
err = self.newFromCloudWafGroup(ctx, userCred, added[i])
err = manager.newFromCloudWafGroup(ctx, userCred, added[i])
if err != nil {
result.AddError(err)
continue
@@ -264,6 +286,20 @@ func (self *SSkuResourcesMeta) SyncWafGroups(ctx context.Context, userCred mccli
return result
}
func fetchCloudEnvs() ([]string, error) {
accounts := []SCloudaccount{}
q := CloudaccountManager.Query("provider", "access_url").In("provider", CloudproviderManager.GetPublicProviderProvidersQuery()).Distinct()
err := q.All(&accounts)
if err != nil {
return nil, errors.Wrapf(err, "q.All")
}
ret := []string{}
for i := range accounts {
ret = append(ret, api.GetCloudEnv(accounts[i].Provider, accounts[i].AccessUrl))
}
return ret, nil
}
func SyncWafGroups(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
err := func() error {
cloudEnvs, err := fetchCloudEnvs()
@@ -271,12 +307,12 @@ func SyncWafGroups(ctx context.Context, userCred mcclient.TokenCredential, isSta
return errors.Wrapf(err, "fetchCloudEnvs")
}
meta, err := FetchSkuResourcesMeta()
meta, err := yunionmeta.FetchYunionmeta(ctx)
if err != nil {
return errors.Wrapf(err, "FetchSkuResourcesMeta")
return errors.Wrapf(err, "FetchYunionmeta")
}
index, err := meta.getWafIndex()
index, err := meta.Index(WafRuleManager.Keyword())
if err != nil {
return errors.Wrapf(err, "getWafIndex")
}
@@ -288,19 +324,13 @@ func SyncWafGroups(ctx context.Context, userCred mcclient.TokenCredential, isSta
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[cloudEnv]
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
if !ok || newMd5 == yunionmeta.EMPTY_MD5 || len(oldMd5) > 0 && newMd5 == oldMd5 {
continue
}
if newMd5 == EMPTY_MD5 {
log.Debugf("%s Waf group is empty skip syncing", cloudEnv)
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Debugf("%s Waf group not Changed skip syncing", cloudEnv)
continue
}
result := meta.SyncWafGroups(ctx, userCred, cloudEnv, isStart)
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
result := WafRuleGroupManager.SyncWafGroups(ctx, userCred, cloudEnv, isStart)
log.Debugf("sync %s waf group result: %s", cloudEnv, result.Result())
}
return nil
+2 -2
View File
@@ -181,8 +181,8 @@ func StartService() {
cron.AddJobAtIntervalsWithStartRun("SyncManagedWafGroups", time.Duration(opts.ServerSkuSyncIntervalMinutes)*time.Minute, models.SyncWafGroups, true)
cron.AddJobEveryFewDays("SyncDBInstanceSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncDBInstanceSkus, true)
cron.AddJobEveryFewDays("SyncNatSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncNatSkus, false)
cron.AddJobEveryFewDays("SyncNasSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncNasSkus, false)
cron.AddJobEveryFewDays("SyncNatSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncNatSkus, true)
cron.AddJobEveryFewDays("SyncNasSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncNasSkus, true)
cron.AddJobEveryFewDays("SyncElasticCacheSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncElasticCacheSkus, true)
cron.AddJobEveryFewDays("StorageSnapshotsRecycle", 1, 2, 0, 0, models.StorageManager.StorageSnapshotsRecycle, false)
@@ -86,13 +86,8 @@ func (self *CloudAccountSyncSkusTask) OnInit(ctx context.Context, obj db.IStanda
}
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, xor bool) compare.SyncResult
type SyncFunc func(ctx context.Context, userCred mcclient.TokenCredential, region *models.SCloudregion, xor bool) compare.SyncResult
var syncFunc SyncFunc
for _, region := range regions {
switch res {
@@ -103,15 +98,15 @@ func (self *CloudAccountSyncSkusTask) OnInit(ctx context.Context, obj db.IStanda
case models.DBInstanceSkuManager.Keyword():
syncFunc = models.DBInstanceSkuManager.SyncDBInstanceSkus
case models.NatSkuManager.Keyword():
result := region.SyncNatSkus(ctx, self.GetUserCred(), meta, false)
result := region.SyncNatSkus(ctx, self.GetUserCred(), false)
log.Infof("Sync %s %s skus for region %s result: %s", region.Provider, res, region.Name, result.Result())
case models.NasSkuManager.Keyword():
result := region.SyncNasSkus(ctx, self.GetUserCred(), meta, false)
result := region.SyncNasSkus(ctx, self.GetUserCred(), false)
log.Infof("Sync %s %s skus for region %s result: %s", region.Provider, res, region.Name, result.Result())
}
if syncFunc != nil {
result := syncFunc(ctx, self.GetUserCred(), &region, meta, false)
result := syncFunc(ctx, self.GetUserCred(), &region, false)
log.Infof("Sync %s %s skus for region %s result: %s", region.Provider, res, region.Name, result.Result())
}
}
@@ -45,13 +45,8 @@ func (self *CloudRegionSyncSkusTask) taskFailed(ctx context.Context, region *mod
func (self *CloudRegionSyncSkusTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) {
region := obj.(*models.SCloudregion)
res, _ := self.GetParams().GetString("resource")
meta, err := models.FetchSkuResourcesMeta()
if err != nil {
self.taskFailed(ctx, region, err.Error())
return
}
type SyncFunc func(ctx context.Context, userCred mcclient.TokenCredential, region *models.SCloudregion, extSkuMeta *models.SSkuResourcesMeta, xor bool) compare.SyncResult
type SyncFunc func(ctx context.Context, userCred mcclient.TokenCredential, region *models.SCloudregion, xor bool) compare.SyncResult
var syncFunc SyncFunc
switch res {
case models.ServerSkuManager.Keyword():
@@ -61,12 +56,12 @@ func (self *CloudRegionSyncSkusTask) OnInit(ctx context.Context, obj db.IStandal
case models.DBInstanceSkuManager.Keyword():
syncFunc = models.DBInstanceSkuManager.SyncDBInstanceSkus
case models.NatSkuManager.Keyword():
result := region.SyncNatSkus(ctx, self.GetUserCred(), meta, false)
result := region.SyncNatSkus(ctx, self.GetUserCred(), false)
log.Infof("Sync %s %s skus for region %s result: %s", region.Provider, res, region.Name, result.Result())
}
if syncFunc != nil {
result := syncFunc(ctx, self.GetUserCred(), region, meta, false)
result := syncFunc(ctx, self.GetUserCred(), region, false)
log.Infof("Sync %s %s skus for region %s result: %s", region.Provider, res, region.Name, result.Result())
if result.IsError() {
self.taskFailed(ctx, region, result.Result())
+1
View File
@@ -0,0 +1 @@
package yunionmeta // import "yunion.io/x/onecloud/pkg/util/yunionmeta"
+181
View File
@@ -0,0 +1,181 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package yunionmeta
import (
"context"
"fmt"
"net/http"
"strings"
"time"
"yunion.io/x/cloudmux/pkg/cloudprovider"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/gotypes"
"yunion.io/x/pkg/util/httputils"
"yunion.io/x/pkg/util/version"
"yunion.io/x/onecloud/pkg/esxi/options"
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/mcclient/modules/compute"
)
const (
EMPTY_MD5 = "d751713988987e9331980363e24189ce"
)
var meta *SSkuResourcesMeta
type SSkuResourcesMeta struct {
// RDS套餐
DBInstanceBase string `json:"dbinstance_base"`
// 虚拟机套餐
ServerBase string `json:"server_base"`
// Redis套餐
ElasticCacheBase string `json:"elastic_cache_base"`
// 公有云镜像
ImageBase string `json:"image_base"`
NatBase string `json:"nat_base"`
NasBase string `json:"nas_base"`
WafBase string `json:"waf_base"`
CloudpolicyBase string `json:"cloudpolicy_base"`
RateBase string `json:"rate_base"`
// 3天过期, 重新刷新
expire time.Time
}
func GetZoneIdBySuffix(zoneMaps map[string]string, suffix string) string {
for externalId, id := range zoneMaps {
if strings.HasSuffix(externalId, suffix) {
return id
}
}
return ""
}
func FetchYunionmeta(ctx context.Context) (*SSkuResourcesMeta, error) {
if !gotypes.IsNil(meta) && meta.expire.After(time.Now()) {
return meta, nil
}
s := auth.GetAdminSession(ctx, "")
transport := httputils.GetTransport(true)
transport.Proxy = options.Options.HttpTransportProxyFunc()
client := &http.Client{Transport: transport}
resp, err := compute.OfflineCloudmeta.GetSkuSourcesMeta(s, client)
if err != nil {
return nil, errors.Wrap(err, "fetchSkuSourceUrls.GetSkuSourcesMeta")
}
meta = &SSkuResourcesMeta{}
err = resp.Unmarshal(meta)
if err != nil {
return nil, errors.Wrap(err, "fetchSkuSourceUrls.Unmarshal")
}
meta.expire = time.Now().AddDate(0, 0, 3)
return meta, nil
}
func (self *SSkuResourcesMeta) request(url string) (jsonutils.JSONObject, error) {
client := httputils.GetAdaptiveTimeoutClient()
header := http.Header{}
header.Set("User-Agent", "vendor/yunion-OneCloud@"+version.Get().GitVersion)
_, resp, err := httputils.JSONRequest(client, context.TODO(), httputils.GET, url, header, nil, false)
return resp, err
}
func (self *SSkuResourcesMeta) _get(url string) ([]jsonutils.JSONObject, error) {
objs, err := self.request(url)
if err != nil {
return nil, errors.Wrapf(err, "request %s", url)
}
var ret []jsonutils.JSONObject
return ret, objs.Unmarshal(&ret)
}
func (self *SSkuResourcesMeta) Get(url string, retVal interface{}) error {
obj, err := self.request(url)
if err != nil {
return errors.Wrapf(err, "request %s", url)
}
return obj.Unmarshal(retVal)
}
func (self *SSkuResourcesMeta) Index(resType string) (map[string]string, error) {
var url string
switch resType {
case "dbinstance_sku":
url = fmt.Sprintf("%s/index.json", self.DBInstanceBase)
case "serversku":
url = fmt.Sprintf("%s/index.json", self.ServerBase)
case "elasticcachesku":
url = fmt.Sprintf("%s/index.json", self.ElasticCacheBase)
case "cloudimage":
url = fmt.Sprintf("%s/index.json", self.ImageBase)
case "nat_sku":
url = fmt.Sprintf("%s/index.json", self.NatBase)
case "nas_sku":
url = fmt.Sprintf("%s/index.json", self.NasBase)
case "waf_rule":
url = fmt.Sprintf("%s/index.json", self.WafBase)
case "cloudpolicy":
url = fmt.Sprintf("%s/index.json", self.CloudpolicyBase)
case "cloudrate":
url = fmt.Sprintf("%s/index.json", self.RateBase)
default:
return nil, errors.Wrapf(cloudprovider.ErrNotFound, resType)
}
ret := map[string]string{}
resp, err := self.request(url)
if err != nil {
return nil, errors.Wrapf(err, "url")
}
return ret, resp.Unmarshal(ret)
}
func (self *SSkuResourcesMeta) List(resType string, regionId string, retVal interface{}) error {
var url string
switch resType {
case "dbinstance_sku":
url = fmt.Sprintf("%s/%s.status.json", self.DBInstanceBase, regionId)
case "serversku":
url = fmt.Sprintf("%s/%s.status.json", self.ServerBase, regionId)
case "elasticcachesku":
url = fmt.Sprintf("%s/%s.status.json", self.ElasticCacheBase, regionId)
case "cloudimage":
url = fmt.Sprintf("%s/%s.status.json", self.ImageBase, regionId)
case "nat_sku":
url = fmt.Sprintf("%s/%s.status.json", self.NatBase, regionId)
case "nas_sku":
url = fmt.Sprintf("%s/%s.status.json", self.NasBase, regionId)
case "waf_rule":
url = fmt.Sprintf("%s/%s.json", self.WafBase, regionId)
case "cloudpolicy":
url = fmt.Sprintf("%s/%s.json", self.CloudpolicyBase, regionId)
default:
return errors.Wrapf(cloudprovider.ErrNotFound, resType)
}
resp, err := self._get(url)
if err != nil {
return errors.Wrapf(err, resType)
}
return jsonutils.Update(retVal, resp)
}