fix(region): optimized qcloud tag sync

This commit is contained in:
Qu Xuan
2021-09-09 17:33:30 +08:00
parent 1af88b6a56
commit 18aff683f9
10 changed files with 91 additions and 99 deletions

View File

@@ -17,6 +17,10 @@ package cloudprovider
import (
"context"
"reflect"
"time"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
)
@@ -40,5 +44,26 @@ func SetTags(ctx context.Context, res ICloudResource, managerId string, tags map
lockman.LockRawObject(ctx, SET_TAGS, managerId)
defer lockman.ReleaseRawObject(ctx, SET_TAGS, managerId)
return res.SetTags(tags, replace)
err := res.SetTags(tags, replace)
if err != nil {
return errors.Wrapf(err, "SetTags")
}
// 避免设置标签后未及时生效,导致本地同步和云上不一致
Wait(time.Second*5, time.Minute, func() (bool, error) {
res.Refresh()
newTags, err := res.GetTags()
if err != nil {
return false, errors.Wrapf(err, "GetTags")
}
for k, v := range tags {
_, ok := newTags[k]
if !ok {
log.Warningf("tag %s:%s not found waitting....", k, v)
return false, nil
}
}
return true, nil
})
return nil
}

View File

@@ -474,6 +474,14 @@ func (lb *SLoadbalancer) GetIRegion() (cloudprovider.ICloudRegion, error) {
return provider.GetIRegionById(region.ExternalId)
}
func (lb *SLoadbalancer) GetILoadbalancer() (cloudprovider.ICloudLoadbalancer, error) {
iRegion, err := lb.GetIRegion()
if err != nil {
return nil, errors.Wrapf(err, "GetIRegion")
}
return iRegion.GetILoadBalancerById(lb.ExternalId)
}
func (lb *SLoadbalancer) GetCreateLoadbalancerParams(iRegion cloudprovider.ICloudRegion) (*cloudprovider.SLoadbalancer, error) {
params := &cloudprovider.SLoadbalancer{
Name: lb.Name,

View File

@@ -303,16 +303,12 @@ func (self *SManagedVirtualizationRegionDriver) RequestStopLoadbalancer(ctx cont
func (self *SManagedVirtualizationRegionDriver) RequestSyncstatusLoadbalancer(ctx context.Context, userCred mcclient.TokenCredential, lb *models.SLoadbalancer, task taskman.ITask) error {
taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) {
iRegion, err := lb.GetIRegion()
iLb, err := lb.GetILoadbalancer()
if err != nil {
return nil, err
return nil, errors.Wrapf(err, "GetILoadbalancer")
}
iLoadbalancer, err := iRegion.GetILoadBalancerById(lb.ExternalId)
if err != nil {
return nil, err
}
models.SyncVirtualResourceMetadata(ctx, userCred, lb, iLoadbalancer)
status := iLoadbalancer.GetStatus()
models.SyncVirtualResourceMetadata(ctx, userCred, lb, iLb)
status := iLb.GetStatus()
if utils.IsInStringArray(status, []string{api.LB_STATUS_ENABLED, api.LB_STATUS_DISABLED}) {
return nil, lb.SetStatus(userCred, status, "")
}
@@ -323,15 +319,11 @@ func (self *SManagedVirtualizationRegionDriver) RequestSyncstatusLoadbalancer(ct
func (self *SManagedVirtualizationRegionDriver) RequestRemoteUpdateLoadbalancer(ctx context.Context, userCred mcclient.TokenCredential, lb *models.SLoadbalancer, replaceTags bool, task taskman.ITask) error {
taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) {
iRegion, err := lb.GetIRegion()
iLb, err := lb.GetILoadbalancer()
if err != nil {
return nil, err
return nil, errors.Wrapf(err, "GetILoadbalancer")
}
iLoadbalancer, err := iRegion.GetILoadBalancerById(lb.ExternalId)
if err != nil {
return nil, err
}
oldTags, err := iLoadbalancer.GetTags()
oldTags, err := iLb.GetTags()
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented {
return nil, nil
@@ -343,7 +335,7 @@ func (self *SManagedVirtualizationRegionDriver) RequestRemoteUpdateLoadbalancer(
return nil, errors.Wrapf(err, "lb.GetAllUserMetadata")
}
tagsUpdateInfo := cloudprovider.TagsUpdateInfo{OldTags: oldTags, NewTags: tags}
err = cloudprovider.SetTags(ctx, iLoadbalancer, lb.ManagerId, tags, replaceTags)
err = cloudprovider.SetTags(ctx, iLb, lb.ManagerId, tags, replaceTags)
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented {
return nil, nil

View File

@@ -18,6 +18,7 @@ import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
@@ -33,22 +34,23 @@ func init() {
taskman.RegisterTask(LoadbalancerRemoteUpdateTask{})
}
func (self *LoadbalancerRemoteUpdateTask) taskFail(ctx context.Context, lb *models.SLoadbalancer, reason jsonutils.JSONObject) {
lb.SetStatus(self.UserCred, api.LB_UPDATE_TAGS_FAILED, reason.String())
self.SetStageFailed(ctx, reason)
func (self *LoadbalancerRemoteUpdateTask) taskFail(ctx context.Context, lb *models.SLoadbalancer, err error) {
lb.SetStatus(self.UserCred, api.LB_UPDATE_TAGS_FAILED, err.Error())
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *LoadbalancerRemoteUpdateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
lb := obj.(*models.SLoadbalancer)
region, err := lb.GetRegion()
if err != nil {
self.taskFail(ctx, lb, jsonutils.NewString(err.Error()))
self.taskFail(ctx, lb, errors.Wrapf(err, "GetRegion"))
return
}
self.SetStage("OnRemoteUpdateComplete", nil)
replaceTags := jsonutils.QueryBoolean(self.Params, "replace_tags", false)
if err := region.GetDriver().RequestRemoteUpdateLoadbalancer(ctx, self.GetUserCred(), lb, replaceTags, self); err != nil {
self.taskFail(ctx, lb, jsonutils.NewString(err.Error()))
self.taskFail(ctx, lb, errors.Wrapf(err, "RequestRemoteUpdateLoadbalancer"))
return
}
}
@@ -58,7 +60,7 @@ func (self *LoadbalancerRemoteUpdateTask) OnRemoteUpdateComplete(ctx context.Con
}
func (self *LoadbalancerRemoteUpdateTask) OnRemoteUpdateCompleteFailed(ctx context.Context, lb *models.SLoadbalancer, data jsonutils.JSONObject) {
self.taskFail(ctx, lb, data)
self.taskFail(ctx, lb, errors.Errorf(data.String()))
}
func (self *LoadbalancerRemoteUpdateTask) OnSyncStatusComplete(ctx context.Context, lb *models.SLoadbalancer, data jsonutils.JSONObject) {

View File

@@ -18,6 +18,7 @@ import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
@@ -35,24 +36,25 @@ func init() {
taskman.RegisterTask(LoadbalancerSyncstatusTask{})
}
func (self *LoadbalancerSyncstatusTask) taskFail(ctx context.Context, lb *models.SLoadbalancer, reason jsonutils.JSONObject) {
lb.SetStatus(self.GetUserCred(), api.LB_STATUS_UNKNOWN, reason.String())
db.OpsLog.LogEvent(lb, db.ACT_SYNC_STATUS, reason, self.UserCred)
logclient.AddActionLogWithStartable(self, lb, logclient.ACT_SYNC_STATUS, reason, self.UserCred, false)
notifyclient.NotifySystemErrorWithCtx(ctx, lb.Id, lb.Name, api.LB_SYNC_CONF_FAILED, reason.String())
self.SetStageFailed(ctx, reason)
func (self *LoadbalancerSyncstatusTask) taskFail(ctx context.Context, lb *models.SLoadbalancer, err error) {
lb.SetStatus(self.GetUserCred(), api.LB_STATUS_UNKNOWN, err.Error())
db.OpsLog.LogEvent(lb, db.ACT_SYNC_STATUS, err, self.UserCred)
logclient.AddActionLogWithStartable(self, lb, logclient.ACT_SYNC_STATUS, err, self.UserCred, false)
notifyclient.NotifySystemErrorWithCtx(ctx, lb.Id, lb.Name, api.LB_SYNC_CONF_FAILED, err.Error())
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *LoadbalancerSyncstatusTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
lb := obj.(*models.SLoadbalancer)
region, err := lb.GetRegion()
if err != nil {
self.taskFail(ctx, lb, jsonutils.NewString(err.Error()))
self.taskFail(ctx, lb, errors.Wrapf(err, "lb.GetRegion"))
return
}
self.SetStage("OnLoadbalancerSyncstatusComplete", nil)
if err := region.GetDriver().RequestSyncstatusLoadbalancer(ctx, self.GetUserCred(), lb, self); err != nil {
self.taskFail(ctx, lb, jsonutils.NewString(err.Error()))
self.taskFail(ctx, lb, errors.Wrapf(err, "RequestSyncstatusLoadbalancer"))
return
}
}
@@ -63,5 +65,5 @@ func (self *LoadbalancerSyncstatusTask) OnLoadbalancerSyncstatusComplete(ctx con
}
func (self *LoadbalancerSyncstatusTask) OnLoadbalancerSyncstatusCompleteFailed(ctx context.Context, lb *models.SLoadbalancer, reason jsonutils.JSONObject) {
self.taskFail(ctx, lb, reason)
self.taskFail(ctx, lb, errors.Errorf(reason.String()))
}

View File

@@ -101,6 +101,7 @@ type Tag struct {
type SInstance struct {
multicloud.SInstanceBase
multicloud.QcloudTags
host *SHost
@@ -167,37 +168,6 @@ func (self *SInstance) GetSecurityGroupIds() ([]string, error) {
return self.SecurityGroupIds, nil
}
func (self *SInstance) GetSysTags() map[string]string {
data := map[string]string{}
if self.image == nil {
image, err := self.host.zone.region.GetImage(self.ImageId)
if err == nil {
self.image = image
}
}
if self.image != nil {
data["os_distribution"] = self.image.OsName
}
priceKey := fmt.Sprintf("%s::%s", self.host.zone.Zone, self.InstanceType)
data["price_key"] = priceKey
data["zone_ext_id"] = self.host.zone.GetGlobalId()
return data
}
func (self *SInstance) GetTags() (map[string]string, error) {
mtags, err := self.host.zone.region.FetchResourceTags("cvm", "instance", []string{self.InstanceId})
if err != nil {
return nil, errors.Wrap(err, "self.host.zone.region.FetchResourceTags")
}
if tags, ok := mtags[self.InstanceId]; ok {
return *tags, nil
}
return map[string]string{}, nil
}
func (self *SInstance) getCloudMetadata() (map[string]string, error) {
mtags, err := self.host.zone.region.FetchResourceTags("cvm", "instance", []string{self.InstanceId})
if err != nil {

View File

@@ -53,6 +53,7 @@ todo:
// https://cloud.tencent.com/document/api/214/30694#LoadBalancer
type SLoadbalancer struct {
multicloud.SLoadbalancerBase
multicloud.QcloudTags
region *SRegion
Status int64 `json:"Status"` // 0创建中1正常运行
@@ -272,24 +273,6 @@ func (self *SLoadbalancer) IsEmulated() bool {
return false
}
func (self *SLoadbalancer) GetTags() (map[string]string, error) {
tags, err := self.region.FetchResourceTags("clb", "clb", []string{self.GetId()})
if err != nil {
return nil, errors.Wrapf(err, "FetchResourceTags")
}
ret := map[string]string{}
if _, ok := tags[self.GetId()]; !ok {
return ret, nil
}
resourceTag := tags[self.GetId()]
if resourceTag != nil {
for k, v := range *resourceTag {
ret[k] = v
}
}
return ret, nil
}
func (self *SLoadbalancer) GetSysTags() map[string]string {
meta := map[string]string{}
meta["Forward"] = strconv.FormatInt(int64(self.Forward), 10)

View File

@@ -898,17 +898,6 @@ func (self *SMySQLInstance) CreateIBackup(opts *cloudprovider.SDBInstanceBackupC
return self.region.CreateMySQLBackup(self.InstanceId, tables)
}
func (self *SMySQLInstance) GetTags() (map[string]string, error) {
tags, err := self.region.FetchResourceTags("cdb", "instanceId", []string{self.GetId()})
if err != nil {
return nil, errors.Wrap(err, "self.region.FetchResourceTags")
}
if _, ok := tags[self.GetId()]; !ok {
return map[string]string{}, nil
}
return *tags[self.GetId()], nil
}
func (self *SMySQLInstance) SetTags(tags map[string]string, replace bool) error {
return self.region.SetResourceTags("cdb", "instanceId", []string{self.InstanceId}, tags, replace)
}

View File

@@ -22,6 +22,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
)
const (
@@ -264,7 +265,7 @@ func (region *SRegion) tagsExist(tags map[string]string) (map[string]string, map
existTags := make(map[string]string)
for i := range tagRows {
tagkv := tagRows[i]
if v, ok := tags[tagkv.TagKey]; ok && tagkv.TagValue == v {
if v, ok := tags[tagkv.TagKey]; ok && (tagkv.TagValue == v || (utils.IsInStringArray(tagkv.TagValue, []string{"", "null"}) && len(v) == 0)) {
// exist
existTags[tagkv.TagKey] = tagkv.TagValue
}
@@ -279,13 +280,27 @@ func (region *SRegion) tagsExist(tags map[string]string) (map[string]string, map
}
func (region *SRegion) SetResourceTags(serviceType, resoureType string, resIds []string, tags map[string]string, replace bool) error {
_, notExist, err := region.tagsExist(tags)
allTags, notExist, err := region.tagsExist(tags)
if err != nil {
return errors.Wrapf(err, "tagsExist")
}
var getTagValue = func(k string) string {
if len(tags[k]) > 0 {
return tags[k]
}
if v, ok := allTags[k]; ok {
return v
}
return "null"
}
for k, v := range notExist {
err := region.createTag(k, v)
err := region.createTag(k, getTagValue(k))
if err != nil {
return errors.Wrapf(err, "createTag %s %s", k, v)
}
}
oldTags, err := region.FetchResourceTags(serviceType, resoureType, resIds)
if err != nil {
return errors.Wrap(err, "FetchTags")
@@ -344,15 +359,15 @@ func (region *SRegion) SetResourceTags(serviceType, resoureType string, resIds [
}
}
for k, ids := range modKeyIds {
err := region.modifyTag(serviceType, resoureType, ids, k, tags[k])
err := region.modifyTag(serviceType, resoureType, ids, k, getTagValue(k))
if err != nil {
return errors.Wrapf(err, "modifyTag %s %s fail %s", k, tags[k], err)
return errors.Wrapf(err, "modifyTag %s %s fail %s", k, getTagValue(k), err)
}
}
for k, ids := range addKeyIds {
err := region.attachTag(serviceType, resoureType, ids, k, tags[k])
err := region.attachTag(serviceType, resoureType, ids, k, getTagValue(k))
if err != nil {
return errors.Wrapf(err, "addTag %s %s fail %s", k, tags[k], err)
return errors.Wrapf(err, "addTag %s %s fail %s", k, getTagValue(k), err)
}
}
return nil

View File

@@ -46,9 +46,15 @@ type QcloudTags struct {
func (self *QcloudTags) GetTags() (map[string]string, error) {
ret := map[string]string{}
for _, tag := range self.TagSet {
if tag.Value == "null" {
tag.Value = ""
}
ret[tag.Key] = tag.Value
}
for _, tag := range self.InstanceTags {
if tag.TagValue == "null" {
tag.TagValue = ""
}
ret[tag.TagKey] = tag.TagValue
}
return ret, nil