diff --git a/pkg/cloudprovider/resourcetags.go b/pkg/cloudprovider/resourcetags.go index 81de0809c8..53dd8456bb 100644 --- a/pkg/cloudprovider/resourcetags.go +++ b/pkg/cloudprovider/resourcetags.go @@ -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 } diff --git a/pkg/compute/models/loadbalancers.go b/pkg/compute/models/loadbalancers.go index c04d50a6b9..0067c94609 100644 --- a/pkg/compute/models/loadbalancers.go +++ b/pkg/compute/models/loadbalancers.go @@ -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, diff --git a/pkg/compute/regiondrivers/managedvirtual.go b/pkg/compute/regiondrivers/managedvirtual.go index d2b06e60d9..df096d65d3 100644 --- a/pkg/compute/regiondrivers/managedvirtual.go +++ b/pkg/compute/regiondrivers/managedvirtual.go @@ -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 diff --git a/pkg/compute/tasks/loadbalancer_remote_update_task.go b/pkg/compute/tasks/loadbalancer_remote_update_task.go index fc209c1fa7..17706b1fa4 100644 --- a/pkg/compute/tasks/loadbalancer_remote_update_task.go +++ b/pkg/compute/tasks/loadbalancer_remote_update_task.go @@ -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) { diff --git a/pkg/compute/tasks/loadbalancer_syncstatus_task.go b/pkg/compute/tasks/loadbalancer_syncstatus_task.go index 7d5eed9476..12d9d15b13 100644 --- a/pkg/compute/tasks/loadbalancer_syncstatus_task.go +++ b/pkg/compute/tasks/loadbalancer_syncstatus_task.go @@ -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())) } diff --git a/pkg/multicloud/qcloud/instance.go b/pkg/multicloud/qcloud/instance.go index eb4f07e1c5..5f50d3d2a3 100644 --- a/pkg/multicloud/qcloud/instance.go +++ b/pkg/multicloud/qcloud/instance.go @@ -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 { diff --git a/pkg/multicloud/qcloud/loadbalancer.go b/pkg/multicloud/qcloud/loadbalancer.go index fec993da45..93b1968f0c 100644 --- a/pkg/multicloud/qcloud/loadbalancer.go +++ b/pkg/multicloud/qcloud/loadbalancer.go @@ -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) diff --git a/pkg/multicloud/qcloud/rds_mysql.go b/pkg/multicloud/qcloud/rds_mysql.go index b104816349..51e6dbdea1 100644 --- a/pkg/multicloud/qcloud/rds_mysql.go +++ b/pkg/multicloud/qcloud/rds_mysql.go @@ -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) } diff --git a/pkg/multicloud/qcloud/tags.go b/pkg/multicloud/qcloud/tags.go index ab8b3d0cf1..187e9a6446 100644 --- a/pkg/multicloud/qcloud/tags.go +++ b/pkg/multicloud/qcloud/tags.go @@ -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 diff --git a/pkg/multicloud/tag_base.go b/pkg/multicloud/tag_base.go index d53c60dbff..34d58bfcd5 100644 --- a/pkg/multicloud/tag_base.go +++ b/pkg/multicloud/tag_base.go @@ -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