From 8408ac8383af9ec732ed4650bcc625476cea7586 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B1=88=E8=BD=A9?= Date: Sat, 13 May 2023 07:51:03 +0800 Subject: [PATCH] fix(region): optimized sku sync (#17029) --- pkg/cloudcommon/db/metadata.go | 1 + pkg/cloudid/models/cloudaccount.go | 29 ++++++++-- pkg/compute/models/cloudimages.go | 9 ++- pkg/compute/models/dbinstance_skus.go | 30 +++++++++- pkg/compute/models/nas_skus.go | 34 +++++++++++- pkg/compute/models/nat_skus.go | 30 +++++++++- pkg/compute/models/skus_tools.go | 70 +++++++++++++++++++----- pkg/compute/models/waf_rule_groups.go | 21 +++---- pkg/mcclient/modules/compute/mod_skus.go | 23 ++++++++ 9 files changed, 209 insertions(+), 38 deletions(-) diff --git a/pkg/cloudcommon/db/metadata.go b/pkg/cloudcommon/db/metadata.go index 8fd9d4523f..56dbc7fbb8 100644 --- a/pkg/cloudcommon/db/metadata.go +++ b/pkg/cloudcommon/db/metadata.go @@ -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" diff --git a/pkg/cloudid/models/cloudaccount.go b/pkg/cloudid/models/cloudaccount.go index 3396f4ad43..19690c3fe5 100644 --- a/pkg/cloudid/models/cloudaccount.go +++ b/pkg/cloudid/models/cloudaccount.go @@ -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) diff --git a/pkg/compute/models/cloudimages.go b/pkg/compute/models/cloudimages.go index b441504e8a..5246fa648b 100644 --- a/pkg/compute/models/cloudimages.go +++ b/pkg/compute/models/cloudimages.go @@ -83,10 +83,15 @@ func SyncPublicCloudImages(ctx context.Context, userCred mcclient.TokenCredentia for i := range regions { region := ®ions[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 { diff --git a/pkg/compute/models/dbinstance_skus.go b/pkg/compute/models/dbinstance_skus.go index dc471b6aa8..71c270381f 100644 --- a/pkg/compute/models/dbinstance_skus.go +++ b/pkg/compute/models/dbinstance_skus.go @@ -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, ®ion, 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) } } diff --git a/pkg/compute/models/nas_skus.go b/pkg/compute/models/nas_skus.go index 06ab8e21d2..6e9b4da548 100644 --- a/pkg/compute/models/nas_skus.go +++ b/pkg/compute/models/nas_skus.go @@ -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 } diff --git a/pkg/compute/models/nat_skus.go b/pkg/compute/models/nat_skus.go index 0911a9387d..7faab05e2b 100644 --- a/pkg/compute/models/nat_skus.go +++ b/pkg/compute/models/nat_skus.go @@ -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) diff --git a/pkg/compute/models/skus_tools.go b/pkg/compute/models/skus_tools.go index 845605de7b..77ccbd8fb1 100644 --- a/pkg/compute/models/skus_tools.go +++ b/pkg/compute/models/skus_tools.go @@ -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 diff --git a/pkg/compute/models/waf_rule_groups.go b/pkg/compute/models/waf_rule_groups.go index 7dcb5e0163..98e885a17e 100644 --- a/pkg/compute/models/waf_rule_groups.go +++ b/pkg/compute/models/waf_rule_groups.go @@ -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 }() diff --git a/pkg/mcclient/modules/compute/mod_skus.go b/pkg/mcclient/modules/compute/mod_skus.go index 83b67f2d43..acc71c0f40 100644 --- a/pkg/mcclient/modules/compute/mod_skus.go +++ b/pkg/mcclient/modules/compute/mod_skus.go @@ -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 +}