fix(region): support eip sync tags (#17796)

This commit is contained in:
屈轩
2023-08-18 15:05:12 +08:00
committed by GitHub
parent f798d8a2b9
commit 6e0df5cf59
13 changed files with 295 additions and 228 deletions
+1 -1
View File
@@ -83,7 +83,7 @@ require (
k8s.io/client-go v0.19.3
k8s.io/cluster-bootstrap v0.19.3
moul.io/http2curl/v2 v2.3.0
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230810103459-4b696c18b138
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230818015659-3fcc9359d8bc
yunion.io/x/executor v0.0.0-20230705125604-c5ac3141db32
yunion.io/x/jsonutils v1.0.1-0.20230613121553-0f3b41e2ef19
yunion.io/x/log v1.0.1-0.20230411060016-feb3f46ab361
+2 -2
View File
@@ -1174,8 +1174,8 @@ sigs.k8s.io/structured-merge-diff/v4 v4.0.1/go.mod h1:bJZC9H9iH24zzfZ/41RGcq60oK
sigs.k8s.io/yaml v1.1.0/go.mod h1:UJmg0vDUVViEyp3mgSv9WPwZCDxu4rQW1olrI1uml+o=
sigs.k8s.io/yaml v1.2.0 h1:kr/MCeFWJWTwyaHoR9c8EjH9OumOmoF9YGiZd7lFm/Q=
sigs.k8s.io/yaml v1.2.0/go.mod h1:yfXDCHCao9+ENCvLSE62v9VSji2MKu5jeNfTrofGhJc=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230810103459-4b696c18b138 h1:AiQTxZuRYQsxP/PVo2pmiUKmgef/0bvE9NkCmYmGZwg=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230810103459-4b696c18b138/go.mod h1:2sgCN7nRPQL3woLfdgqLDd92vwAHqtlz3KKiHxC5BAw=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230818015659-3fcc9359d8bc h1:drAGYF6oevqZaCTuccQjN/fMxoEB+oK07bEsyvduKvE=
yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230818015659-3fcc9359d8bc/go.mod h1:2sgCN7nRPQL3woLfdgqLDd92vwAHqtlz3KKiHxC5BAw=
yunion.io/x/executor v0.0.0-20230705125604-c5ac3141db32 h1:v7POYkQwo1XzOxBoIoRVr/k0V9Y5JyjpshlIFa9raug=
yunion.io/x/executor v0.0.0-20230705125604-c5ac3141db32/go.mod h1:Uxuou9WQIeJXNpy7t2fPLL0BYLvLiMvGQwY7Qc6aSws=
yunion.io/x/jsonutils v0.0.0-20190625054549-a964e1e8a051/go.mod h1:4N0/RVzsYL3kH3WE/H1BjUQdFiWu50JGCFQuuy+Z634=
+5
View File
@@ -143,3 +143,8 @@ type ElasticDissociateInput struct {
// default: false
AutoDelete bool `json:"auto_delete"`
}
type ElasticipRemoteUpdateInput struct {
// 是否覆盖替换所有标签
ReplaceTags *bool `json:"replace_tags" help:"replace all remote tags"`
}
+29
View File
@@ -1937,6 +1937,35 @@ func (eip *SElasticip) GetUsages() []db.IUsage {
}
}
func (self *SElasticip) PerformRemoteUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ElasticipRemoteUpdateInput) (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 *SElasticip) 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, "EipRemoteUpdateTask", 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 *SElasticip) 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)
}
}
func (manager *SElasticipManager) ListItemExportKeys(ctx context.Context,
q *sqlchemy.SQuery,
userCred mcclient.TokenCredential,
@@ -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/cloudmux/pkg/cloudprovider"
"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/compute/models"
"yunion.io/x/onecloud/pkg/util/logclient"
)
type EipRemoteUpdateTask struct {
taskman.STask
}
func init() {
taskman.RegisterTask(EipRemoteUpdateTask{})
}
func (self *EipRemoteUpdateTask) taskFail(ctx context.Context, eip *models.SElasticip, err error) {
eip.SetStatus(self.UserCred, apis.STATUS_UPDATE_TAGS_FAILED, err.Error())
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
func (self *EipRemoteUpdateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
eip := obj.(*models.SElasticip)
replaceTags := jsonutils.QueryBoolean(self.Params, "replace_tags", false)
iEip, err := eip.GetIEip(ctx)
if err != nil {
self.taskFail(ctx, eip, errors.Wrapf(err, "GetIEip"))
return
}
oldTags, err := iEip.GetTags()
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented {
self.OnRemoteUpdateComplete(ctx, eip, nil)
return
}
self.taskFail(ctx, eip, errors.Wrapf(err, "GetTags"))
return
}
tags, err := eip.GetAllUserMetadata()
if err != nil {
self.taskFail(ctx, eip, errors.Wrapf(err, "GetAllUserMetadata"))
return
}
tagsUpdateInfo := cloudprovider.TagsUpdateInfo{OldTags: oldTags, NewTags: tags}
err = cloudprovider.SetTags(ctx, iEip, eip.ManagerId, tags, replaceTags)
if err != nil {
if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented {
self.OnRemoteUpdateComplete(ctx, eip, nil)
return
}
logclient.AddActionLogWithStartable(self, eip, logclient.ACT_UPDATE_TAGS, err, self.GetUserCred(), false)
self.taskFail(ctx, eip, err)
return
}
logclient.AddActionLogWithStartable(self, eip, logclient.ACT_UPDATE_TAGS, tagsUpdateInfo, self.GetUserCred(), true)
self.OnRemoteUpdateComplete(ctx, eip, nil)
}
func (self *EipRemoteUpdateTask) OnRemoteUpdateComplete(ctx context.Context, eip *models.SElasticip, data jsonutils.JSONObject) {
self.SetStage("OnSyncStatusComplete", nil)
models.StartResourceSyncStatusTask(ctx, self.UserCred, eip, "EipSyncstatusTask", self.GetTaskId())
}
func (self *EipRemoteUpdateTask) OnSyncStatusComplete(ctx context.Context, eip *models.SElasticip, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
func (self *EipRemoteUpdateTask) OnSyncStatusCompleteFailed(ctx context.Context, eip *models.SElasticip, data jsonutils.JSONObject) {
self.SetStageFailed(ctx, data)
}
@@ -79,7 +79,7 @@ func (self *NatGatewayRemoteUpdateTask) OnInit(ctx context.Context, obj db.IStan
return
}
logclient.AddActionLogWithStartable(self, nat, logclient.ACT_UPDATE_TAGS, err, self.GetUserCred(), false)
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
self.taskFail(ctx, nat, err)
return
}
logclient.AddActionLogWithStartable(self, nat, logclient.ACT_UPDATE_TAGS, tagsUpdateInfo, self.GetUserCred(), true)
+1 -1
View File
@@ -1438,7 +1438,7 @@ sigs.k8s.io/structured-merge-diff/v4/value
# sigs.k8s.io/yaml v1.2.0
## explicit; go 1.12
sigs.k8s.io/yaml
# yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230810103459-4b696c18b138
# yunion.io/x/cloudmux v0.3.10-0-alpha.1.0.20230818015659-3fcc9359d8bc
## explicit; go 1.18
yunion.io/x/cloudmux/pkg/apis
yunion.io/x/cloudmux/pkg/apis/billing
+4
View File
@@ -294,3 +294,7 @@ func (self *SEipAddress) GetExpiredAt() time.Time {
func (self *SEipAddress) GetProjectId() string {
return ""
}
func (self *SEipAddress) SetTags(tags map[string]string, replace bool) error {
return self.region.setTags("elastic-ip", self.AllocationId, tags, replace)
}
+2 -68
View File
@@ -457,13 +457,7 @@ func (self *SInstance) DeleteVM(ctx context.Context) error {
}
func (self *SInstance) UpdateVM(ctx context.Context, name string) error {
addTags := map[string]string{}
addTags["Name"] = name
Arn := self.GetArn()
err := self.host.zone.region.TagResources([]string{Arn}, addTags)
if err != nil {
return errors.Wrapf(err, "self.host.zone.region.TagResources([]string{%s}, %s)", Arn, jsonutils.Marshal(addTags).String())
}
self.SetTags(map[string]string{"Name": name}, false)
return nil
}
@@ -1272,54 +1266,7 @@ func (self *SInstance) GetError() error {
}
func (self *SInstance) SetTags(tags map[string]string, replace bool) error {
ec2Client, err := self.host.zone.region.getEc2Client()
if err != nil {
return errors.Wrap(err, "getEc2Client")
}
oldTagsJson, err := FetchTags(ec2Client, self.InstanceId)
if err != nil {
return errors.Wrapf(err, "FetchTags(self.host.zone.region.ec2Client, %s)", self.InstanceId)
}
_oldTags := TagSpec{}
err = oldTagsJson.Unmarshal(&_oldTags.Tags)
if err != nil {
return errors.Wrapf(err, "(%s).Unmarshal(oldTags)", oldTagsJson.String())
}
oldTags, err := _oldTags.GetTags()
if err != nil {
return errors.Wrapf(err, "GetTags")
}
addTags := map[string]string{}
for k, v := range tags {
if _, ok := oldTags[k]; !ok {
addTags[k] = v
} else {
if oldTags[k] != v {
addTags[k] = v
}
}
}
delTags := []string{}
if replace {
for k := range oldTags {
if _, ok := tags[k]; !ok {
if !strings.HasPrefix(k, "aws:") && k != "Name" {
delTags = append(delTags, k)
}
}
}
}
Arn := self.GetArn()
err = self.host.zone.region.UntagResources([]string{Arn}, delTags)
if err != nil {
return errors.Wrapf(err, "UntagResources")
}
delete(addTags, "Name")
err = self.host.zone.region.TagResources([]string{Arn}, addTags)
if err != nil {
return errors.Wrapf(err, "TagResources")
}
return nil
return self.host.zone.region.setTags("instance", self.InstanceId, tags, replace)
}
func (self *SInstance) GetAccountId() string {
@@ -1331,19 +1278,6 @@ func (self *SInstance) GetAccountId() string {
return identity.Account
}
func (self *SInstance) GetArn() string {
partition := ""
switch self.host.zone.region.client.GetAccessEnv() {
case api.CLOUD_ACCESS_ENV_AWS_GLOBAL:
partition = "aws"
case api.CLOUD_ACCESS_ENV_AWS_CHINA:
partition = "aws-cn"
default:
partition = "aws"
}
return fmt.Sprintf("arn:%s:ec2:%s:%s:instance/%s", partition, self.host.zone.region.GetId(), self.GetAccountId(), self.InstanceId)
}
func (self *SRegion) SaveImage(instanceId string, opts *cloudprovider.SaveImageOptions) (*SImage, error) {
params := map[string]string{
"Description": opts.Notes,
+66 -7
View File
@@ -121,7 +121,7 @@ func (self *SElb) GetSysTags() map[string]string {
}
func (self *SElb) GetTags() (map[string]string, error) {
return self.region.FetchElbTags(self.LoadBalancerArn)
return self.region.DescribeElbTags(self.LoadBalancerArn)
}
func (self *SElb) GetAddress() string {
@@ -391,18 +391,77 @@ func (self *SRegion) CreateElbBackendgroup(opts *cloudprovider.SLoadbalancerBack
}
func (self *SElb) SetTags(tags map[string]string, replace bool) error {
oldTags, err := self.region.FetchElbTags(self.LoadBalancerArn)
return self.region.setElbTags(self.LoadBalancerArn, tags, replace)
}
func (self *SRegion) setElbTags(arn string, tags map[string]string, replace bool) error {
oldTags, err := self.DescribeElbTags(arn)
if err != nil {
return errors.Wrapf(err, "self.region.FetchElbTags(%s)", self.LoadBalancerArn)
return errors.Wrapf(err, "DescribeElbTags")
}
err = self.region.UpdateResourceTags(self.LoadBalancerArn, oldTags, tags, replace)
if err != nil {
return errors.Wrap(err, "self.region.UpdateResourceTags(self.LoadBalancerArn, oldTags, tags, replace)")
added, removed := map[string]string{}, map[string]string{}
for k, v := range tags {
oldValue, ok := oldTags[k]
if !ok {
added[k] = v
} else if oldValue != v {
removed[k] = oldValue
added[k] = v
}
}
if replace {
for k, v := range oldTags {
newValue, ok := tags[k]
if !ok {
removed[k] = v
} else if v != newValue {
added[k] = newValue
removed[k] = v
}
}
}
if len(removed) > 0 {
err = self.RemoveElbTags(arn, removed)
if err != nil {
return errors.Wrapf(err, "RemoveElbTags %s", removed)
}
}
if len(added) > 0 {
return self.AddElbTags(arn, added)
}
return nil
}
func (self *SRegion) FetchElbTags(arn string) (map[string]string, error) {
func (self *SRegion) AddElbTags(arn string, tags map[string]string) error {
params := map[string]string{
"ResourceArns.member.1": arn,
}
idx := 1
for k, v := range tags {
params[fmt.Sprintf("Tags.member.%d.Key", idx)] = k
params[fmt.Sprintf("Tags.member.%d.Value", idx)] = v
idx++
}
ret := struct {
}{}
return self.elbRequest("AddTags", params, &ret)
}
func (self *SRegion) RemoveElbTags(arn string, tags map[string]string) error {
params := map[string]string{
"ResourceArns.member.1": arn,
}
idx := 1
for k := range tags {
params[fmt.Sprintf("TagKeys.member.%d", idx)] = k
idx++
}
ret := struct {
}{}
return self.elbRequest("RemoveTags", params, &ret)
}
func (self *SRegion) DescribeElbTags(arn string) (map[string]string, error) {
ret := struct {
TagDescriptions []struct {
ResourceArn string `xml:"ResourceArn"`
+90
View File
@@ -242,3 +242,93 @@ func (self *SVpc) GetINatGateways() ([]cloudprovider.ICloudNatGateway, error) {
}
return ret, nil
}
func (self *SRegion) DescribeTags(resType, instanceId string) (map[string]string, error) {
params := map[string]string{
"Filter.1.Name": "resource-id",
"Filter.1.Value.1": instanceId,
"Filter.2.Name": "resource-type",
"Filter.2.Value.1": resType,
}
ret := struct {
NextToken string
AwsTags
}{}
err := self.ec2Request("DescribeTags", params, &ret)
if err != nil {
return nil, err
}
return ret.GetTags()
}
func (self *SRegion) DeleteTags(instanceId string, tags map[string]string) error {
params := map[string]string{
"ResourceId.1": instanceId,
}
idx := 1
for k, v := range tags {
params[fmt.Sprintf("Tag.%d.Key", idx)] = k
params[fmt.Sprintf("Tag.%d.Value", idx)] = v
idx++
}
ret := struct {
}{}
return self.ec2Request("DeleteTags", params, &ret)
}
func (self *SRegion) CreateTags(instanceId string, tags map[string]string) error {
params := map[string]string{
"ResourceId.1": instanceId,
}
idx := 1
for k, v := range tags {
params[fmt.Sprintf("Tag.%d.Key", idx)] = k
params[fmt.Sprintf("Tag.%d.Value", idx)] = v
idx++
}
ret := struct {
}{}
return self.ec2Request("CreateTags", params, &ret)
}
func (self *SRegion) setTags(resType, resId string, tags map[string]string, replace bool) error {
oldTags, err := self.DescribeTags(resType, resId)
if err != nil {
return errors.Wrapf(err, "DescribeTags")
}
added, removed := map[string]string{}, map[string]string{}
for k, v := range tags {
oldValue, ok := oldTags[k]
if !ok {
added[k] = v
} else if oldValue != v {
removed[k] = oldValue
added[k] = v
}
}
if replace {
for k, v := range oldTags {
newValue, ok := tags[k]
if !ok {
removed[k] = v
} else if v != newValue {
added[k] = newValue
removed[k] = v
}
}
}
if len(removed) > 0 {
err = self.DeleteTags(resId, removed)
if err != nil {
return errors.Wrapf(err, "DeleteTags %s", removed)
}
}
if len(added) > 0 {
return self.CreateTags(resId, added)
}
return nil
}
func (self *SNatGateway) SetTags(tags map[string]string, replace bool) error {
return self.region.setTags("natgateway", self.NatGatewayId, tags, replace)
}
-121
View File
@@ -1,121 +0,0 @@
// 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 aws
import (
"strings"
"github.com/aws/aws-sdk-go/service/resourcegroupstaggingapi"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/cloudmux/pkg/cloudprovider"
)
func (self *SRegion) TagResources(arns []string, tags map[string]string) error {
if len(tags) == 0 {
return nil
}
client, err := self.getResourceGroupTagClient()
if err != nil {
return errors.Wrap(err, "self.getResourceGroupTagClient()")
}
params := resourcegroupstaggingapi.TagResourcesInput{}
arnInput := []*string{}
for i := range arns {
arnInput = append(arnInput, &arns[i])
}
tagsInput := make(map[string]*string)
tagValues := []string{}
for k, v := range tags {
tagValues = append(tagValues, v)
tagsInput[k] = &tagValues[len(tagValues)-1]
}
params.SetResourceARNList(arnInput)
params.SetTags(tagsInput)
out, err := client.TagResources(&params)
if err != nil {
return errors.Wrapf(err, "client.TagResources(%s)", jsonutils.Marshal(params).String())
}
if out != nil && len(out.FailedResourcesMap) > 0 {
return errors.Wrapf(cloudprovider.ErrNotSupported, "client.TagResources(%s),error:%s", jsonutils.Marshal(params).String(), jsonutils.Marshal(out).String())
}
return nil
}
func (self *SRegion) UntagResources(arns []string, tagKeys []string) error {
if len(tagKeys) == 0 {
return nil
}
client, err := self.getResourceGroupTagClient()
if err != nil {
return errors.Wrap(err, "self.getResourceGroupTagClient()")
}
params := resourcegroupstaggingapi.UntagResourcesInput{}
arnInput := []*string{}
for i := range arns {
arnInput = append(arnInput, &arns[i])
}
delTagKeysInput := []*string{}
for i := range tagKeys {
delTagKeysInput = append(delTagKeysInput, &tagKeys[i])
}
params.SetResourceARNList(arnInput)
params.SetTagKeys(delTagKeysInput)
out, err := client.UntagResources(&params)
if err != nil {
return errors.Wrapf(err, "client.UntagResources(%s)", jsonutils.Marshal(params).String())
}
if out != nil && len(out.FailedResourcesMap) > 0 {
return errors.Wrapf(cloudprovider.ErrNotSupported, "client.UntagResources(%s),error:%s", jsonutils.Marshal(params).String(), jsonutils.Marshal(out).String())
}
return nil
}
func (self *SRegion) UpdateResourceTags(arn string, oldTags, tags map[string]string, replace bool) error {
addTags := map[string]string{}
for k, v := range tags {
if strings.HasPrefix(k, "aws:") {
return errors.Wrap(cloudprovider.ErrNotSupported, "The aws: prefix is reserved for AWS use")
}
if _, ok := oldTags[k]; !ok {
addTags[k] = v
} else {
if oldTags[k] != v {
addTags[k] = v
}
}
}
delTags := []string{}
if replace {
for k := range oldTags {
if _, ok := tags[k]; !ok {
if !strings.HasPrefix(k, "aws:") {
delTags = append(delTags, k)
}
}
}
}
err := self.UntagResources([]string{arn}, delTags)
if err != nil {
return errors.Wrapf(err, "self.host.zone.region.UntagResources([]string{%s}, %s)", arn, jsonutils.Marshal(delTags).String())
}
err = self.TagResources([]string{arn}, addTags)
if err != nil {
return errors.Wrapf(err, "self.host.zone.region.TagResources([]string{%s}, %s)", arn, jsonutils.Marshal(addTags).String())
}
return nil
}
-27
View File
@@ -21,7 +21,6 @@ import (
"github.com/aws/aws-sdk-go/service/ec2"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/cloudmux/pkg/cloudprovider"
@@ -270,32 +269,6 @@ func NextDeviceName(curDeviceNames []string) (string, error) {
return "", fmt.Errorf("disk devicename out of index, current deivces: %s", currents)
}
// fetch tags
func FetchTags(client *ec2.EC2, resourceId string) (*jsonutils.JSONDict, error) {
result := jsonutils.NewDict()
params := &ec2.DescribeTagsInput{}
filters := []*ec2.Filter{}
if len(resourceId) == 0 {
return result, fmt.Errorf("resource id should not be empty")
}
// todo: add resource type filter
filters = AppendSingleValueFilter(filters, "resource-id", resourceId)
params.SetFilters(filters)
ret, err := client.DescribeTags(params)
if err != nil {
return result, err
}
for _, tag := range ret.Tags {
if tag.Key != nil && tag.Value != nil {
result.Set(*tag.Key, jsonutils.NewString(*tag.Value))
}
}
return result, nil
}
// error
func parseNotFoundError(err error) error {
if err == nil {