From 9dea6d2a7d03832f89f3151fb0c1b26fa6b14968 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B1=88=E8=BD=A9?= Date: Wed, 26 Oct 2022 02:10:38 +0800 Subject: [PATCH] fix(region): skip disk sync by system tags (#15231) --- pkg/compute/models/disks.go | 24 ++++++++++++++--- pkg/compute/models/hosts.go | 44 ++++++++++++++++---------------- pkg/compute/options/options.go | 4 +-- pkg/multicloud/aliyun/storage.go | 2 ++ pkg/multicloud/tag_base.go | 6 ++--- 5 files changed, 49 insertions(+), 31 deletions(-) diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index f3f44fe8e0..07de46887d 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -1476,17 +1476,33 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To } for i := 0; i < len(commondb); i += 1 { + skip, key := IsNeedSkipSync(commonext[i]) + if skip { + log.Infof("delete disk %s(%s) with tag key: %s", commonext[i].GetName(), commonext[i].GetGlobalId(), key) + err := commondb[i].purge(ctx, userCred) + if err != nil { + syncResult.DeleteError(err) + continue + } + syncResult.Delete() + continue + } err = commondb[i].syncWithCloudDisk(ctx, userCred, provider, commonext[i], -1, syncOwnerId, storage.ManagerId) if err != nil { syncResult.UpdateError(err) - } else { - localDisks = append(localDisks, commondb[i]) - remoteDisks = append(remoteDisks, commonext[i]) - syncResult.Update() + continue } + localDisks = append(localDisks, commondb[i]) + remoteDisks = append(remoteDisks, commonext[i]) + syncResult.Update() } for i := 0; i < len(added); i += 1 { + skip, key := IsNeedSkipSync(added[i]) + if skip { + log.Infof("skip disk %s(%s) sync with tag key: %s", added[i].GetName(), added[i].GetGlobalId(), key) + continue + } extId := added[i].GetGlobalId() _disk, err := db.FetchByExternalIdAndManagerId(manager, extId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery { sq := StorageManager.Query().SubQuery() diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index f26c59ca99..fed0b0659a 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -2433,6 +2433,26 @@ type SGuestSyncResult struct { IsNew bool } +func IsNeedSkipSync(ext cloudprovider.ICloudResource) (bool, string) { + if len(options.Options.SkipServerBySysTagKeys) == 0 && len(options.Options.SkipServerBySysTagKeys) == 0 { + return false, "" + } + keys := strings.Split(options.Options.SkipServerBySysTagKeys, ",") + for key := range ext.GetSysTags() { + if utils.IsInStringArray(key, keys) { + return true, key + } + } + userKeys := strings.Split(options.Options.SkipServerByUserTagKeys, ",") + tags, _ := ext.GetTags() + for key := range tags { + if utils.IsInStringArray(key, userKeys) { + return true, key + } + } + return false, "" +} + func (self *SHost) SyncHostVMs(ctx context.Context, userCred mcclient.TokenCredential, iprovider cloudprovider.ICloudProvider, vms []cloudprovider.ICloudVM, syncOwnerId mcclient.IIdentityProvider) ([]SGuestSyncResult, compare.SyncResult) { lockman.LockRawObject(ctx, "guests", self.Id) defer lockman.ReleaseRawObject(ctx, "guests", self.Id) @@ -2464,26 +2484,6 @@ func (self *SHost) SyncHostVMs(ctx context.Context, userCred mcclient.TokenCrede return nil, syncResult } - skipFunc := func(ext cloudprovider.ICloudVM) (bool, string) { - if len(options.Options.SkipServerBySysTagKeys) == 0 && len(options.Options.SkipServerBySysTagKeys) == 0 { - return false, "" - } - keys := strings.Split(options.Options.SkipServerBySysTagKeys, ",") - for key := range ext.GetSysTags() { - if utils.IsInStringArray(key, keys) { - return true, key - } - } - userKeys := strings.Split(options.Options.SkipServerByUserTagKeys, ",") - tags, _ := ext.GetTags() - for key := range tags { - if utils.IsInStringArray(key, userKeys) { - return true, key - } - } - return false, "" - } - for i := 0; i < len(removed); i += 1 { err := removed[i].syncRemoveCloudVM(ctx, userCred) if err != nil { @@ -2494,7 +2494,7 @@ func (self *SHost) SyncHostVMs(ctx context.Context, userCred mcclient.TokenCrede } for i := 0; i < len(commondb); i += 1 { - skip, key := skipFunc(commonext[i]) + skip, key := IsNeedSkipSync(commonext[i]) if skip { log.Infof("delete server %s(%s) with system tag key: %s", commonext[i].GetName(), commonext[i].GetGlobalId(), key) err := commondb[i].purge(ctx, userCred) @@ -2520,7 +2520,7 @@ func (self *SHost) SyncHostVMs(ctx context.Context, userCred mcclient.TokenCrede } for i := 0; i < len(added); i += 1 { - skip, key := skipFunc(added[i]) + skip, key := IsNeedSkipSync(added[i]) if skip { log.Infof("skip server %s(%s) sync with system tag key: %s", added[i].GetName(), added[i].GetGlobalId(), key) continue diff --git a/pkg/compute/options/options.go b/pkg/compute/options/options.go index b8c5c0c6c5..d6fd300c7a 100644 --- a/pkg/compute/options/options.go +++ b/pkg/compute/options/options.go @@ -187,8 +187,8 @@ type ComputeOptions struct { DefaultIPAllocationDirection string `help:"default IP allocation direction" default:"stepdown"` // 弹性伸缩中的ecs一般会有特殊的系统标签,通过指定这些标签可以忽略这部分ecs的同步, 指定多个key需要以 ',' 分隔 - SkipServerBySysTagKeys string `help:"skip server sync and create with system tags" default:"acs:autoscaling:scalingGroupId"` - SkipServerByUserTagKeys string `help:"skip server sync and create with user tags" default:""` + SkipServerBySysTagKeys string `help:"skip server,disk sync and create with system tags" default:"acs:autoscaling:scalingGroupId"` + SkipServerByUserTagKeys string `help:"skip server,disk sync and create with user tags" default:""` EnableAwsMonitorAgent bool `help:"enable aws monitor agent" default:"true"` diff --git a/pkg/multicloud/aliyun/storage.go b/pkg/multicloud/aliyun/storage.go index c80810ccf5..ab7de429ec 100644 --- a/pkg/multicloud/aliyun/storage.go +++ b/pkg/multicloud/aliyun/storage.go @@ -82,6 +82,8 @@ func (self *SStorage) GetIDisks() ([]cloudprovider.ICloudDisk, error) { } performanceLevel := "" switch self.storageType { + case api.STORAGE_CLOUD_ESSD: + performanceLevel = "PL1" case api.STORAGE_CLOUD_ESSD_PL2: performanceLevel = "PL2" case api.STORAGE_CLOUD_ESSD_PL3: diff --git a/pkg/multicloud/tag_base.go b/pkg/multicloud/tag_base.go index 2c05319e46..c46dc43e98 100644 --- a/pkg/multicloud/tag_base.go +++ b/pkg/multicloud/tag_base.go @@ -141,8 +141,8 @@ func (self *AliyunTags) GetTags() (map[string]string, error) { ret := map[string]string{} for _, tag := range self.Tags.Tag { if strings.HasPrefix(tag.TagKey, "aliyun") || strings.HasPrefix(tag.TagKey, "acs:") || - strings.HasSuffix(tag.Key, "aliyun") || strings.HasPrefix(tag.Key, "acs:") || - strings.HasSuffix(tag.Key, "ack.") { // k8s + strings.HasPrefix(tag.Key, "aliyun") || strings.HasPrefix(tag.Key, "acs:") || + strings.HasPrefix(tag.Key, "ack.") || strings.HasPrefix(tag.TagKey, "ack.") { // k8s continue } if len(tag.TagKey) > 0 { @@ -160,7 +160,7 @@ func (self *AliyunTags) GetSysTags() map[string]string { for _, tag := range self.Tags.Tag { if strings.HasPrefix(tag.TagKey, "aliyun") || strings.HasPrefix(tag.TagKey, "acs:") || strings.HasPrefix(tag.Key, "aliyun") || strings.HasPrefix(tag.Key, "acs:") || - strings.HasPrefix(tag.Key, "ack.") { // k8s + strings.HasPrefix(tag.Key, "ack.") || strings.HasPrefix(tag.TagKey, "ack.") { // k8s if len(tag.TagKey) > 0 { ret[tag.TagKey] = tag.TagValue } else if len(tag.Key) > 0 {