fix(region): optimized sku sync (#17029)

This commit is contained in:
屈轩
2023-05-13 07:51:03 +08:00
committed by GitHub
parent 24e04ba84c
commit 8408ac8383
9 changed files with 209 additions and 38 deletions
+1
View File
@@ -46,6 +46,7 @@ const (
USER_TAG_PREFIX = dbapi.USER_TAG_PREFIX
SYS_CLOUD_TAG_PREFIX = dbapi.SYS_CLOUD_TAG_PREFIX
CLASS_TAG_PREFIX = dbapi.CLASS_TAT_PREFIX
SKU_METADAT_KEY = "md5"
// TAG_DELETE_RANGE_USER = "user"
// TAG_DELETE_RANGE_CLOUD = CLOUD_TAG_PREFIX // "cloud"
+24 -5
View File
@@ -850,6 +850,10 @@ 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 +868,28 @@ func (self *SCloudaccount) SyncSystemCloudpoliciesFromCloud(ctx context.Context,
transport := httputils.GetTransport(true)
transport.Proxy = options.Options.HttpTransportProxyFunc()
client := &http.Client{Transport: transport}
meta, err := modules.OfflineCloudmeta.GetSkuSourcesMeta(s, client)
policyBase, index, err := modules.OfflineCloudmeta.GetSkuIndex(s, client, "cloudpolicy_base")
if err != nil {
return errors.Wrap(err, "GetSkuSourcesMeta")
return errors.Wrapf(err, "get cloudpolicy")
}
policyBase, err := meta.GetString("cloudpolicy_base")
if err != nil {
return errors.Wrapf(err, "missing policy base url")
skuMeta := &SCloudpolicy{}
skuMeta.SetModelManager(CloudpolicyManager, skuMeta)
skuMeta.Id = self.Provider
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)
return nil
}
policyUrl := strings.TrimSuffix(policyBase, "/") + fmt.Sprintf("/%s.json", self.Provider)
+7 -2
View File
@@ -83,10 +83,15 @@ func SyncPublicCloudImages(ctx context.Context, userCred mcclient.TokenCredentia
for i := range regions {
region := &regions[i]
oldMd5, _ := imageIndex[region.ExternalId]
skuMeta := &SCloudimage{}
skuMeta.SetModelManager(CloudimageManager, skuMeta)
skuMeta.Id = region.ExternalId
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[region.ExternalId]
if ok {
imageIndex[region.ExternalId] = newMd5
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
+28 -2
View File
@@ -603,15 +603,41 @@ func SyncRegionDBInstanceSkus(ctx context.Context, userCred mcclient.TokenCreden
return
}
index, err := meta.getSkuIndex(meta.DBInstanceBase)
if err != nil {
log.Errorf("get rds sku index error: %v", err)
return
}
for _, region := range cloudregions {
if !region.GetDriver().IsSupportedDBInstance() {
log.Infof("region %s(%s) not support dbinstance, skip sync", region.Name, region.Id)
log.Debugf("region %s(%s) not support dbinstance, skip sync", region.Name, region.Id)
continue
}
skuMeta := &SDBInstanceSku{}
skuMeta.SetModelManager(DBInstanceSkuManager, 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 DBInstance skus is empty skip syncing", region.Name)
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Debugf("%s DBInstance skus not changed skip syncing", region.Name)
continue
}
result := DBInstanceSkuManager.SyncDBInstanceSkus(ctx, userCred, &region, meta, xor)
msg := result.Result()
notes := fmt.Sprintf("sync rds sku for region %s result: %s", region.Name, msg)
log.Infof(notes)
log.Debugf(notes)
}
}
+31 -3
View File
@@ -314,15 +314,43 @@ func SyncRegionNasSkus(ctx context.Context, userCred mcclient.TokenCredential, r
return errors.Wrapf(err, "FetchSkuResourcesMeta")
}
index, err := meta.getSkuIndex(meta.NasBase)
if err != nil {
log.Errorf("get nas sku index error: %v", err)
return err
}
for i := range regions {
if !regions[i].GetDriver().IsSupportedNas() {
log.Infof("region %s(%s) not support nas, skip sync", regions[i].Name, regions[i].Id)
region := regions[i]
if !region.GetDriver().IsSupportedNas() {
log.Debugf("region %s(%s) not support nas, skip sync", regions[i].Name, regions[i].Id)
continue
}
skuMeta := &SNasSku{}
skuMeta.SetModelManager(NasSkuManager, 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 nas skus is empty skip syncing", region.Name)
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Debugf("%s nas skus not changed skip syncing", region.Name)
continue
}
result := regions[i].SyncNasSkus(ctx, userCred, meta, xor)
msg := result.Result()
notes := fmt.Sprintf("SyncNasSkus for region %s result: %s", regions[i].Name, msg)
log.Infof(notes)
log.Debugf(notes)
}
return nil
}
+29 -1
View File
@@ -317,11 +317,39 @@ func SyncRegionNatSkus(ctx context.Context, userCred mcclient.TokenCredential, r
return errors.Wrapf(err, "FetchSkuResourcesMeta")
}
index, err := meta.getSkuIndex(meta.NatBase)
if err != nil {
log.Errorf("get nat sku index error: %v", err)
return err
}
for i := range regions {
if !regions[i].GetDriver().IsSupportedNatGateway() {
region := regions[i]
if !region.GetDriver().IsSupportedNatGateway() {
log.Infof("region %s(%s) not support nat, skip sync", regions[i].Name, regions[i].Id)
continue
}
skuMeta := &SNatSku{}
skuMeta.SetModelManager(NatSkuManager, 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.Infof("%s Nat Skus is empty skip syncing", region.Name)
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Infof("%s Nat Skus not Changed skip syncing", region.Name)
continue
}
result := regions[i].SyncNatSkus(ctx, userCred, meta, xor)
msg := result.Result()
notes := fmt.Sprintf("SyncNatSkus for region %s result: %s", regions[i].Name, msg)
+55 -15
View File
@@ -56,9 +56,6 @@ type SSkuResourcesMeta struct {
WafBase string `json:"waf_base"`
}
var skuIndex = map[string]string{}
var imageIndex = map[string]string{}
func (self *SSkuResourcesMeta) getZoneIdBySuffix(zoneMaps map[string]string, suffix string) string {
for externalId, id := range zoneMaps {
if strings.HasSuffix(externalId, suffix) {
@@ -401,6 +398,19 @@ func (self *SSkuResourcesMeta) getServerSkuIndex() (map[string]string, error) {
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 {
@@ -473,18 +483,43 @@ func SyncElasticCacheSkus(ctx context.Context, userCred mcclient.TokenCredential
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() {
result := ElasticcacheSkuManager.SyncElasticcacheSkus(ctx, userCred, region, meta, false)
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)
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)
}
}
@@ -536,25 +571,30 @@ func SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, isSt
cloudregions := fetchSkuSyncCloudregions()
for i := range cloudregions {
region := &cloudregions[i]
oldMd5, _ := skuIndex[region.ExternalId]
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 {
skuIndex[region.ExternalId] = newMd5
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
if newMd5 == EMPTY_MD5 {
log.Infof("%s Server Skus is empty skip syncing", region.Name)
log.Debugf("%s server skus is empty skip syncing", region.Name)
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Infof("%s Server Skus not Changed skip syncing", region.Name)
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.Infof(notes)
log.Debugf(notes)
}
// 清理无效的sku
+11 -10
View File
@@ -37,8 +37,6 @@ type SWafRuleGroupManager struct {
db.SExternalizedResourceBaseManager
}
var wafIndex map[string]string
var WafRuleGroupManager *SWafRuleGroupManager
func init() {
@@ -50,7 +48,6 @@ func init() {
"waf_rule_groups",
),
}
wafIndex = map[string]string{}
WafRuleGroupManager.SetVirtualObject(WafRuleGroupManager)
}
@@ -285,22 +282,26 @@ func SyncWafGroups(ctx context.Context, userCred mcclient.TokenCredential, isSta
}
for _, cloudEnv := range cloudEnvs {
skuMeta := &SWafRuleGroup{}
skuMeta.SetModelManager(WafRuleGroupManager, skuMeta)
skuMeta.Id = cloudEnv
oldMd5 := db.Metadata.GetStringValue(ctx, skuMeta, db.SKU_METADAT_KEY, userCred)
newMd5, ok := index[cloudEnv]
if !ok {
continue
if ok {
db.Metadata.SetValue(ctx, skuMeta, db.SKU_METADAT_KEY, newMd5, userCred)
}
oldMd5, _ := wafIndex[cloudEnv]
if newMd5 == EMPTY_MD5 {
log.Infof("%s Waf group is empty skip syncing", cloudEnv)
log.Debugf("%s Waf group is empty skip syncing", cloudEnv)
continue
}
if len(oldMd5) > 0 && newMd5 == oldMd5 {
log.Infof("%s Waf group not Changed skip syncing", cloudEnv)
log.Debugf("%s Waf group not Changed skip syncing", cloudEnv)
continue
}
result := meta.SyncWafGroups(ctx, userCred, cloudEnv, isStart)
log.Infof("sync %s waf group result: %s", cloudEnv, result.Result())
wafIndex[cloudEnv] = newMd5
log.Debugf("sync %s waf group result: %s", cloudEnv, result.Result())
}
return nil
}()
+23
View File
@@ -21,6 +21,7 @@ import (
"strings"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/httputils"
"yunion.io/x/pkg/util/printutils"
@@ -101,3 +102,25 @@ func (self *OfflineCloudmetaManager) GetSkuSourcesMeta(s *mcclient.ClientSession
_, body, err := httputils.JSONRequest(client, context.TODO(), "GET", url, nil, nil, false)
return body, err
}
func (self *OfflineCloudmetaManager) GetSkuIndex(s *mcclient.ClientSession, client *http.Client, res string) (string, map[string]string, error) {
meta, err := self.GetSkuSourcesMeta(s, client)
if err != nil {
return "", nil, errors.Wrapf(err, "GetSkuSourcesMeta")
}
base, err := meta.GetString(res)
if err != nil {
return "", nil, errors.Wrapf(err, "get %s", res)
}
url := fmt.Sprintf("%s/index.json", base)
_, body, err := httputils.JSONRequest(client, context.TODO(), "GET", url, nil, nil, false)
if err != nil {
return "", map[string]string{}, errors.Wrapf(err, "request")
}
ret := map[string]string{}
err = body.Unmarshal(ret)
if err != nil {
return "", map[string]string{}, errors.Wrapf(err, "resp.Unmarshal")
}
return base, ret, nil
}