diff --git a/pkg/compute/models/cdn_domains.go b/pkg/compute/models/cdn_domains.go index baf4ab461b..8a1f36df0a 100644 --- a/pkg/compute/models/cdn_domains.go +++ b/pkg/compute/models/cdn_domains.go @@ -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) + } +} diff --git a/pkg/compute/models/natgateways.go b/pkg/compute/models/natgateways.go index 501f0f6095..e6053c8421 100644 --- a/pkg/compute/models/natgateways.go +++ b/pkg/compute/models/natgateways.go @@ -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) + } +} diff --git a/pkg/compute/tasks/cdn_domain_remote_update_task.go b/pkg/compute/tasks/cdn_domain_remote_update_task.go new file mode 100644 index 0000000000..f2cbf09005 --- /dev/null +++ b/pkg/compute/tasks/cdn_domain_remote_update_task.go @@ -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) +} diff --git a/pkg/compute/tasks/natgateway_remote_update_task.go b/pkg/compute/tasks/natgateway_remote_update_task.go new file mode 100644 index 0000000000..5e7739302b --- /dev/null +++ b/pkg/compute/tasks/natgateway_remote_update_task.go @@ -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) +} diff --git a/pkg/compute/tasks/natgateway_syncstatus_task.go b/pkg/compute/tasks/natgateway_syncstatus_task.go index 97a0553bd6..9e8338b692 100644 --- a/pkg/compute/tasks/natgateway_syncstatus_task.go +++ b/pkg/compute/tasks/natgateway_syncstatus_task.go @@ -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) } diff --git a/pkg/multicloud/aliyun/natgateway.go b/pkg/multicloud/aliyun/natgateway.go index 9383dbcdde..8f388f2faa 100644 --- a/pkg/multicloud/aliyun/natgateway.go +++ b/pkg/multicloud/aliyun/natgateway.go @@ -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) +} diff --git a/pkg/multicloud/qcloud/cdn.go b/pkg/multicloud/qcloud/cdn.go index e09e9c2a6d..a112d0e165 100644 --- a/pkg/multicloud/qcloud/cdn.go +++ b/pkg/multicloud/qcloud/cdn.go @@ -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") diff --git a/pkg/multicloud/tag_base.go b/pkg/multicloud/tag_base.go index 118b213120..210899bf78 100644 --- a/pkg/multicloud/tag_base.go +++ b/pkg/multicloud/tag_base.go @@ -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 }