mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
feat(region):sync kafka tag
fmt
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
Reference in New Issue
Block a user