Merge pull request #12398 from ioito/hotfix/qx-nat-cdn-tags

Hotfix/qx nat cdn tags
This commit is contained in:
Zexi Li
2021-10-14 12:45:03 +08:00
committed by GitHub
8 changed files with 338 additions and 15 deletions
+35
View File
@@ -19,10 +19,12 @@ import (
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/compare"
"yunion.io/x/sqlchemy"
"yunion.io/x/onecloud/pkg/apis"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
@@ -389,3 +391,36 @@ func (self *SCDNDomain) PerformSyncstatus(ctx context.Context, userCred mcclient
func (self *SCDNDomain) StartSyncstatus(ctx context.Context, userCred mcclient.TokenCredential, parentTaskId string) error {
return StartResourceSyncStatusTask(ctx, userCred, self, "CDNDomainSyncstatusTask", parentTaskId)
}
func (self *SCDNDomain) AllowPerformRemoteUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return self.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, self, "remote-update")
}
func (self *SCDNDomain) PerformRemoteUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.MongoDBRemoteUpdateInput) (jsonutils.JSONObject, error) {
err := self.StartRemoteUpdateTask(ctx, userCred, (input.ReplaceTags != nil && *input.ReplaceTags), "")
if err != nil {
return nil, errors.Wrap(err, "StartRemoteUpdateTask")
}
return nil, nil
}
func (self *SCDNDomain) StartRemoteUpdateTask(ctx context.Context, userCred mcclient.TokenCredential, replaceTags bool, parentTaskId string) error {
data := jsonutils.NewDict()
data.Add(jsonutils.NewBool(replaceTags), "replace_tags")
task, err := taskman.TaskManager.NewTask(ctx, "CDNDomainRemoteUpdateTask", self, userCred, data, parentTaskId, "", nil)
if err != nil {
return errors.Wrap(err, "NewTask")
}
self.SetStatus(userCred, apis.STATUS_UPDATE_TAGS, "StartRemoteUpdateTask")
return task.ScheduleRun(nil)
}
func (self *SCDNDomain) OnMetadataUpdated(ctx context.Context, userCred mcclient.TokenCredential) {
if len(self.ExternalId) == 0 {
return
}
err := self.StartRemoteUpdateTask(ctx, userCred, true, "")
if err != nil {
log.Errorf("StartRemoteUpdateTask fail: %s", err)
}
}
+35 -2
View File
@@ -426,7 +426,6 @@ func (manager *SNatGatewayManager) SyncNatGateways(ctx context.Context, userCred
syncResult.UpdateError(err)
continue
}
syncMetadata(ctx, userCred, &commondb[i], commonext[i])
localNatGateways = append(localNatGateways, commondb[i])
remoteNatGateways = append(remoteNatGateways, commonext[i])
syncResult.Update()
@@ -438,7 +437,6 @@ func (manager *SNatGatewayManager) SyncNatGateways(ctx context.Context, userCred
syncResult.AddError(err)
continue
}
syncMetadata(ctx, userCred, routeTableNew, added[i])
localNatGateways = append(localNatGateways, *routeTableNew)
remoteNatGateways = append(remoteNatGateways, added[i])
syncResult.Add()
@@ -505,6 +503,7 @@ func (self *SNatGateway) SyncWithCloudNatGateway(ctx context.Context, userCred m
return err
}
syncMetadata(ctx, userCred, self, extNat)
SyncCloudDomain(userCred, self, provider.GetOwnerId())
db.OpsLog.LogSyncUpdate(self, diff, userCred)
@@ -564,6 +563,7 @@ func (manager *SNatGatewayManager) newFromCloudNatGateway(ctx context.Context, u
}
SyncCloudDomain(userCred, &nat, provider.GetOwnerId())
syncMetadata(ctx, userCred, &nat, extNat)
db.OpsLog.LogEvent(&nat, db.ACT_CREATE, nat.GetShortDesc(ctx), userCred)
@@ -1088,3 +1088,36 @@ func (manager *SNatGatewayManager) DeleteExpiredPostpaids(ctx context.Context, u
nats[i].StartNatGatewayDeleteTask(ctx, userCred, nil)
}
}
func (self *SNatGateway) AllowPerformRemoteUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return self.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, self, "remote-update")
}
func (self *SNatGateway) PerformRemoteUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.MongoDBRemoteUpdateInput) (jsonutils.JSONObject, error) {
err := self.StartRemoteUpdateTask(ctx, userCred, (input.ReplaceTags != nil && *input.ReplaceTags), "")
if err != nil {
return nil, errors.Wrap(err, "StartRemoteUpdateTask")
}
return nil, nil
}
func (self *SNatGateway) StartRemoteUpdateTask(ctx context.Context, userCred mcclient.TokenCredential, replaceTags bool, parentTaskId string) error {
data := jsonutils.NewDict()
data.Add(jsonutils.NewBool(replaceTags), "replace_tags")
task, err := taskman.TaskManager.NewTask(ctx, "NatGatewayRemoteUpdateTask", self, userCred, data, parentTaskId, "", nil)
if err != nil {
return errors.Wrap(err, "NewTask")
}
self.SetStatus(userCred, apis.STATUS_UPDATE_TAGS, "StartRemoteUpdateTask")
return task.ScheduleRun(nil)
}
func (self *SNatGateway) OnMetadataUpdated(ctx context.Context, userCred mcclient.TokenCredential) {
if len(self.ExternalId) == 0 {
return
}
err := self.StartRemoteUpdateTask(ctx, userCred, true, "")
if err != nil {
log.Errorf("StartRemoteUpdateTask fail: %s", err)
}
}
@@ -0,0 +1,94 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package tasks
import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/apis"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/util/logclient"
)
type CDNDomainRemoteUpdateTask struct {
taskman.STask
}
func init() {
taskman.RegisterTask(CDNDomainRemoteUpdateTask{})
}
func (self *CDNDomainRemoteUpdateTask) taskFail(ctx context.Context, cdn *models.SCDNDomain, err error) {
cdn.SetStatus(self.UserCred, apis.STATUS_UPDATE_TAGS_FAILED, err.Error())
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *CDNDomainRemoteUpdateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
cdn := obj.(*models.SCDNDomain)
replaceTags := jsonutils.QueryBoolean(self.Params, "replace_tags", false)
iCDNDomain, err := cdn.GetICloudCDNDomain()
if err != nil {
self.taskFail(ctx, cdn, errors.Wrapf(err, "GetICloudCDNDomain"))
return
}
oldTags, err := iCDNDomain.GetTags()
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented {
self.OnRemoteUpdateComplete(ctx, cdn, nil)
return
}
self.taskFail(ctx, cdn, errors.Wrapf(err, "GetTags"))
return
}
tags, err := cdn.GetAllUserMetadata()
if err != nil {
self.taskFail(ctx, cdn, errors.Wrapf(err, "GetAllUserMetadata"))
return
}
tagsUpdateInfo := cloudprovider.TagsUpdateInfo{OldTags: oldTags, NewTags: tags}
err = cloudprovider.SetTags(ctx, iCDNDomain, cdn.ManagerId, tags, replaceTags)
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented {
self.OnRemoteUpdateComplete(ctx, cdn, nil)
return
}
logclient.AddActionLogWithStartable(self, cdn, logclient.ACT_UPDATE_TAGS, err, self.GetUserCred(), false)
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
return
}
logclient.AddActionLogWithStartable(self, cdn, logclient.ACT_UPDATE_TAGS, tagsUpdateInfo, self.GetUserCred(), true)
self.OnRemoteUpdateComplete(ctx, cdn, nil)
}
func (self *CDNDomainRemoteUpdateTask) OnRemoteUpdateComplete(ctx context.Context, cdn *models.SCDNDomain, data jsonutils.JSONObject) {
self.SetStage("OnSyncStatusComplete", nil)
models.StartResourceSyncStatusTask(ctx, self.UserCred, cdn, "CDNDomainSyncstatusTask", self.GetTaskId())
}
func (self *CDNDomainRemoteUpdateTask) OnSyncStatusComplete(ctx context.Context, cdn *models.SCDNDomain, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
func (self *CDNDomainRemoteUpdateTask) OnSyncStatusCompleteFailed(ctx context.Context, cdn *models.SCDNDomain, data jsonutils.JSONObject) {
self.SetStageFailed(ctx, data)
}
@@ -0,0 +1,100 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package tasks
import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/apis"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/util/logclient"
)
type NatGatewayRemoteUpdateTask struct {
taskman.STask
}
func init() {
taskman.RegisterTask(NatGatewayRemoteUpdateTask{})
}
func (self *NatGatewayRemoteUpdateTask) taskFail(ctx context.Context, nat *models.SNatGateway, err error) {
nat.SetStatus(self.UserCred, apis.STATUS_UPDATE_TAGS_FAILED, err.Error())
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *NatGatewayRemoteUpdateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
nat := obj.(*models.SNatGateway)
replaceTags := jsonutils.QueryBoolean(self.Params, "replace_tags", false)
vpc, err := nat.GetVpc()
if err != nil {
self.taskFail(ctx, nat, errors.Wrapf(err, "GetVpc"))
return
}
iNatGateway, err := nat.GetINatGateway()
if err != nil {
self.taskFail(ctx, nat, errors.Wrapf(err, "GetINatGateway"))
return
}
oldTags, err := iNatGateway.GetTags()
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented {
self.OnRemoteUpdateComplete(ctx, nat, nil)
return
}
self.taskFail(ctx, nat, errors.Wrapf(err, "GetTags"))
return
}
tags, err := nat.GetAllUserMetadata()
if err != nil {
self.taskFail(ctx, nat, errors.Wrapf(err, "GetAllUserMetadata"))
return
}
tagsUpdateInfo := cloudprovider.TagsUpdateInfo{OldTags: oldTags, NewTags: tags}
err = cloudprovider.SetTags(ctx, iNatGateway, vpc.ManagerId, tags, replaceTags)
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented {
self.OnRemoteUpdateComplete(ctx, nat, nil)
return
}
logclient.AddActionLogWithStartable(self, nat, logclient.ACT_UPDATE_TAGS, err, self.GetUserCred(), false)
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
return
}
logclient.AddActionLogWithStartable(self, nat, logclient.ACT_UPDATE_TAGS, tagsUpdateInfo, self.GetUserCred(), true)
self.OnRemoteUpdateComplete(ctx, nat, nil)
}
func (self *NatGatewayRemoteUpdateTask) OnRemoteUpdateComplete(ctx context.Context, nat *models.SNatGateway, data jsonutils.JSONObject) {
self.SetStage("OnSyncStatusComplete", nil)
models.StartResourceSyncStatusTask(ctx, self.UserCred, nat, "NatGatewaySyncstatusTask", self.GetTaskId())
}
func (self *NatGatewayRemoteUpdateTask) OnSyncStatusComplete(ctx context.Context, nat *models.SNatGateway, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
func (self *NatGatewayRemoteUpdateTask) OnSyncStatusCompleteFailed(ctx context.Context, nat *models.SNatGateway, data jsonutils.JSONObject) {
self.SetStageFailed(ctx, data)
}
+11 -11
View File
@@ -16,9 +16,9 @@ package tasks
import (
"context"
"fmt"
"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,26 +35,26 @@ func init() {
taskman.RegisterTask(NatGatewaySyncstatusTask{})
}
func (self *NatGatewaySyncstatusTask) taskFailed(ctx context.Context, natgateway *models.SNatGateway, reason jsonutils.JSONObject) {
natgateway.SetStatus(self.GetUserCred(), api.NAT_STATUS_UNKNOWN, reason.String())
self.SetStageFailed(ctx, reason)
func (self *NatGatewaySyncstatusTask) taskFailed(ctx context.Context, natgateway *models.SNatGateway, err error) {
natgateway.SetStatus(self.GetUserCred(), api.NAT_STATUS_UNKNOWN, err.Error())
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
db.OpsLog.LogEvent(natgateway, db.ACT_SYNC_STATUS, natgateway.GetShortDesc(ctx), self.GetUserCred())
logclient.AddActionLogWithContext(ctx, natgateway, logclient.ACT_SYNC_STATUS, reason, self.UserCred, false)
logclient.AddActionLogWithContext(ctx, natgateway, logclient.ACT_SYNC_STATUS, err, self.UserCred, false)
}
func (self *NatGatewaySyncstatusTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
natgateway := obj.(*models.SNatGateway)
region, _ := natgateway.GetRegion()
if region == nil {
self.taskFailed(ctx, natgateway, jsonutils.NewString(fmt.Sprintf("failed to found cloudregion for natgateway %s(%s)", natgateway.Name, natgateway.Id)))
region, err := natgateway.GetRegion()
if err != nil {
self.taskFailed(ctx, natgateway, errors.Wrapf(err, "GetRegion"))
return
}
self.SetStage("OnNatGatewaySyncStatusComplete", nil)
err := region.GetDriver().RequestSyncNatGatewayStatus(ctx, self.GetUserCred(), natgateway, self)
err = region.GetDriver().RequestSyncNatGatewayStatus(ctx, self.GetUserCred(), natgateway, self)
if err != nil {
self.taskFailed(ctx, natgateway, jsonutils.NewString(err.Error()))
self.taskFailed(ctx, natgateway, errors.Wrapf(err, "RequestSyncNatGatewayStatus"))
return
}
}
@@ -64,5 +64,5 @@ func (self *NatGatewaySyncstatusTask) OnNatGatewaySyncStatusComplete(ctx context
}
func (self *NatGatewaySyncstatusTask) OnNatGatewaySyncStatusCompleteFailed(ctx context.Context, natgateway *models.SNatGateway, data jsonutils.JSONObject) {
self.taskFailed(ctx, natgateway, data)
self.SetStageFailed(ctx, data)
}
+29
View File
@@ -16,6 +16,7 @@ package aliyun
import (
"fmt"
"strings"
"time"
"yunion.io/x/jsonutils"
@@ -350,3 +351,31 @@ func (self *SRegion) DeleteNatGateway(natId string, isForce bool) error {
_, err := self.vpcRequest("DeleteNatGateway", params)
return errors.Wrapf(err, "DeleteNatGateway")
}
func (self *SNatGateway) GetTags() (map[string]string, error) {
_, tags, err := self.vpc.region.ListSysAndUserTags(ALIYUN_SERVICE_VPC, "NATGATEWAY", self.NatGatewayId)
if err != nil {
return nil, errors.Wrapf(err, "ListTags")
}
tagMaps := map[string]string{}
for k, v := range tags {
tagMaps[strings.ToLower(k)] = v
}
return tagMaps, nil
}
func (self *SNatGateway) GetSysTags() map[string]string {
tags, _, err := self.vpc.region.ListSysAndUserTags(ALIYUN_SERVICE_VPC, "NATGATEWAY", self.NatGatewayId)
if err != nil {
return nil
}
tagMaps := map[string]string{}
for k, v := range tags {
tagMaps[strings.ToLower(k)] = v
}
return tagMaps
}
func (self *SNatGateway) SetTags(tags map[string]string, replace bool) error {
return self.vpc.region.SetResourceTags(ALIYUN_SERVICE_VPC, "NATGATEWAY", self.GetId(), tags, replace)
}
+25 -2
View File
@@ -18,6 +18,7 @@ import (
"fmt"
"strconv"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/cloudprovider"
@@ -83,10 +84,32 @@ func (self *SCdnDomain) GetServiceType() string {
return self.ServiceType
}
func (self *SCdnDomain) Refresh() error {
domains, _, err := self.client.DescribeCdnDomains([]string{self.Domain}, nil, "", 0, 100)
if err != nil {
return errors.Wrapf(err, "DescribeCdnDomains")
}
if len(domains) == 0 {
return errors.Wrapf(cloudprovider.ErrNotFound, self.Domain)
}
if len(domains) > 1 {
return errors.Wrapf(cloudprovider.ErrDuplicateId, self.Domain)
}
return jsonutils.Update(self, domains[0])
}
func (self *SCdnDomain) Delete() error {
return self.client.DeleteCdnDomain(self.Domain)
}
func (self *SCdnDomain) SetTags(tags map[string]string, replace bool) error {
region, err := self.client.getDefaultRegion()
if err != nil {
return errors.Wrapf(err, "getDefaultRegion")
}
return region.SetResourceTags("cdn", "domain", []string{self.Domain}, tags, replace)
}
func (self *SQcloudClient) DeleteCdnDomain(domain string) error {
params := map[string]string{
"Domain": domain,
@@ -170,9 +193,9 @@ func (client *SQcloudClient) DescribeCdnDomains(domains, origins []string, domai
filterIndex++
}
resp, err := client.cdnRequest("DescribeDomains", params)
resp, err := client.cdnRequest("DescribeDomainsConfig", params)
if err != nil {
return nil, 0, errors.Wrapf(err, "DescribeDomains %s", params)
return nil, 0, errors.Wrapf(err, "DescribeDomainsConfig %s", params)
}
cdnDomains := []SCdnDomain{}
err = resp.Unmarshal(&cdnDomains, "Domains")
+9
View File
@@ -46,6 +46,8 @@ type QcloudTags struct {
TagList []STag
// Kafka
Tags []STag
// Cdn
Tag []STag
}
func (self *QcloudTags) GetTags() (map[string]string, error) {
@@ -74,6 +76,13 @@ func (self *QcloudTags) GetTags() (map[string]string, error) {
}
ret[tag.TagKey] = tag.TagValue
}
for _, tag := range self.Tag {
if tag.TagValue == "null" {
tag.TagValue = ""
}
ret[tag.TagKey] = tag.TagValue
}
return ret, nil
}