From 115bb1290f694b514236f4bff4a5ffcff139aceb Mon Sep 17 00:00:00 2001 From: silence Date: Wed, 28 Dec 2022 16:12:39 +0800 Subject: [PATCH] feat(region):sync kafka tag fmt --- pkg/apis/compute/kafka.go | 2 + pkg/compute/models/kafka.go | 28 ++++++++ pkg/compute/models/regiondrivers.go | 5 ++ pkg/compute/regiondrivers/base.go | 4 ++ pkg/compute/regiondrivers/managedvirtual.go | 32 +++++++++ pkg/compute/tasks/kafka_remote_update_task.go | 72 +++++++++++++++++++ 6 files changed, 143 insertions(+) create mode 100644 pkg/compute/tasks/kafka_remote_update_task.go diff --git a/pkg/apis/compute/kafka.go b/pkg/apis/compute/kafka.go index 30a777b48b..e0426579a3 100644 --- a/pkg/apis/compute/kafka.go +++ b/pkg/apis/compute/kafka.go @@ -27,6 +27,8 @@ const ( KAFKA_STATUS_DELETING = compute.KAFKA_STATUS_DELETING KAFKA_STATUS_DELETE_FAILED = "delete_failed" KAFKA_STATUS_UNKNOWN = compute.KAFKA_STATUS_UNKNOWN + KAFKA_UPDATE_TAGS = "update_tags" + KAFKA_UPDATE_TAGS_FAILED = "update_tags_fail" ) type KafkaCreateInput struct { diff --git a/pkg/compute/models/kafka.go b/pkg/compute/models/kafka.go index a4618e96ae..a9f07e8bc7 100644 --- a/pkg/compute/models/kafka.go +++ b/pkg/compute/models/kafka.go @@ -645,6 +645,10 @@ func (self *SKafka) PerformSyncstatus(ctx context.Context, userCred mcclient.Tok return nil, StartResourceSyncStatusTask(ctx, userCred, self, "KafkaSyncstatusTask", "") } +func (self *SKafka) StartKafkaSyncTask(ctx context.Context, userCred mcclient.TokenCredential, parentTaskId string) error { + return StartResourceSyncStatusTask(ctx, userCred, self, "KafkaSyncstatusTask", parentTaskId) +} + // 获取Kafka Topic列表 func (self *SKafka) GetDetailsTopics(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) ([]cloudprovider.SKafkaTopic, error) { iKafka, err := self.GetIKafka(ctx) @@ -653,3 +657,27 @@ func (self *SKafka) GetDetailsTopics(ctx context.Context, userCred mcclient.Toke } return iKafka.GetTopics() } + +func (self *SKafka) StartRemoteUpdateTask(ctx context.Context, userCred mcclient.TokenCredential, replaceTags bool, parentTaskId string) error { + data := jsonutils.NewDict() + if replaceTags { + data.Add(jsonutils.JSONTrue, "replace_tags") + } + if task, err := taskman.TaskManager.NewTask(ctx, "KafkaRemoteUpdateTask", self, userCred, data, parentTaskId, "", nil); err != nil { + return errors.Wrap(err, "Start ElasticSearchRemoteUpdateTask") + } else { + self.SetStatus(userCred, api.ELASTIC_SEARCH_UPDATE_TAGS, "StartRemoteUpdateTask") + task.ScheduleRun(nil) + } + return nil +} + +func (self *SKafka) 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/regiondrivers.go b/pkg/compute/models/regiondrivers.go index 3909a5b54f..95decada57 100644 --- a/pkg/compute/models/regiondrivers.go +++ b/pkg/compute/models/regiondrivers.go @@ -39,6 +39,7 @@ type IRegionDriver interface { IElasticcacheBackup IDBInstanceDriver IElasticSearchDriver + IKafkaDrive ValidateCreateLoadbalancerData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, data *api.LoadbalancerCreateInput) (*api.LoadbalancerCreateInput, error) RequestCreateLoadbalancerInstance(ctx context.Context, userCred mcclient.TokenCredential, lb *SLoadbalancer, input *api.LoadbalancerCreateInput, task taskman.ITask) error @@ -275,6 +276,10 @@ type IElasticSearchDriver interface { RequestRemoteUpdateElasticSearch(ctx context.Context, userCred mcclient.TokenCredential, elasticcache *SElasticSearch, replaceTags bool, task taskman.ITask) error } +type IKafkaDrive interface { + RequestRemoteUpdateKafka(ctx context.Context, userCred mcclient.TokenCredential, kafka *SKafka, replaceTags bool, task taskman.ITask) error +} + var regionDrivers map[string]IRegionDriver func init() { diff --git a/pkg/compute/regiondrivers/base.go b/pkg/compute/regiondrivers/base.go index 8181098d9b..44891d48de 100644 --- a/pkg/compute/regiondrivers/base.go +++ b/pkg/compute/regiondrivers/base.go @@ -525,3 +525,7 @@ func (self *SBaseRegionDriver) RequestUnpackInstanceBackup(ctx context.Context, func (self *SBaseRegionDriver) RequestRemoteUpdateElasticSearch(ctx context.Context, userCred mcclient.TokenCredential, elasticcache *models.SElasticSearch, replaceTags bool, task taskman.ITask) error { return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestRemoteUpdateElasticSearch") } + +func (self *SBaseRegionDriver) RequestRemoteUpdateKafka(ctx context.Context, userCred mcclient.TokenCredential, kafka *models.SKafka, replaceTags bool, task taskman.ITask) error { + return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestRemoteUpdateElasticSearch") +} diff --git a/pkg/compute/regiondrivers/managedvirtual.go b/pkg/compute/regiondrivers/managedvirtual.go index 68199be207..8e16ba9409 100644 --- a/pkg/compute/regiondrivers/managedvirtual.go +++ b/pkg/compute/regiondrivers/managedvirtual.go @@ -3280,3 +3280,35 @@ func (self *SManagedVirtualizationRegionDriver) RequestRemoteUpdateElasticSearch }) return nil } + +func (self *SManagedVirtualizationRegionDriver) RequestRemoteUpdateKafka(ctx context.Context, userCred mcclient.TokenCredential, instance *models.SKafka, replaceTags bool, task taskman.ITask) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + kafka, err := instance.GetIKafka(ctx) + if err != nil { + return nil, errors.Wrap(err, "instance.GetIKafka") + } + oldTags, err := kafka.GetTags() + if err != nil { + if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented { + return nil, nil + } + return nil, errors.Wrap(err, "ies.GetTags()") + } + tags, err := instance.GetAllUserMetadata() + if err != nil { + return nil, errors.Wrapf(err, "instance.GetAllUserMetadata") + } + tagsUpdateInfo := cloudprovider.TagsUpdateInfo{OldTags: oldTags, NewTags: tags} + err = cloudprovider.SetTags(ctx, kafka, instance.ManagerId, tags, replaceTags) + if err != nil { + if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented { + return nil, nil + } + logclient.AddActionLogWithStartable(task, instance, logclient.ACT_UPDATE_TAGS, err, userCred, false) + return nil, errors.Wrap(err, "ies.SetTags") + } + logclient.AddActionLogWithStartable(task, instance, logclient.ACT_UPDATE_TAGS, tagsUpdateInfo, userCred, true) + return nil, nil + }) + return nil +} diff --git a/pkg/compute/tasks/kafka_remote_update_task.go b/pkg/compute/tasks/kafka_remote_update_task.go new file mode 100644 index 0000000000..4c969e9b60 --- /dev/null +++ b/pkg/compute/tasks/kafka_remote_update_task.go @@ -0,0 +1,72 @@ +// 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" + "fmt" + + "yunion.io/x/jsonutils" + + api "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" +) + +type KafkaRemoteUpdateTask struct { + taskman.STask +} + +func init() { + taskman.RegisterTask(KafkaRemoteUpdateTask{}) +} + +func (self *KafkaRemoteUpdateTask) taskFail(ctx context.Context, kafka *models.SKafka, reason jsonutils.JSONObject) { + kafka.SetStatus(self.UserCred, api.KAFKA_UPDATE_TAGS_FAILED, reason.String()) + self.SetStageFailed(ctx, reason) +} + +func (self *KafkaRemoteUpdateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + kafka := obj.(*models.SKafka) + region, _ := kafka.GetRegion() + if region == nil { + self.taskFail(ctx, kafka, jsonutils.NewString(fmt.Sprintf("failed to find region for elastic search %s", kafka.GetName()))) + return + } + self.SetStage("OnRemoteUpdateComplete", nil) + replaceTags := jsonutils.QueryBoolean(self.Params, "replace_tags", false) + + if err := region.GetDriver().RequestRemoteUpdateKafka(ctx, self.GetUserCred(), kafka, replaceTags, self); err != nil { + self.taskFail(ctx, kafka, jsonutils.NewString(err.Error())) + } +} + +func (self *KafkaRemoteUpdateTask) OnRemoteUpdateComplete(ctx context.Context, kafka *models.SKafka, data jsonutils.JSONObject) { + self.SetStage("OnSyncStatusComplete", nil) + kafka.StartKafkaSyncTask(ctx, self.UserCred, self.GetTaskId()) +} + +func (self *KafkaRemoteUpdateTask) OnRemoteUpdateCompleteFailed(ctx context.Context, kafka *models.SKafka, data jsonutils.JSONObject) { + self.taskFail(ctx, kafka, data) +} + +func (self *KafkaRemoteUpdateTask) OnSyncStatusComplete(ctx context.Context, kafka *models.SKafka, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} + +func (self *KafkaRemoteUpdateTask) OnSyncStatusCompleteFailed(ctx context.Context, kafka *models.SKafka, data jsonutils.JSONObject) { + self.SetStageFailed(ctx, data) +}