diff --git a/pkg/apis/compute/filesystem.go b/pkg/apis/compute/filesystem.go index 32846b37e1..053b36b3f0 100644 --- a/pkg/apis/compute/filesystem.go +++ b/pkg/apis/compute/filesystem.go @@ -36,6 +36,9 @@ const ( // 删除中 NAS_STATUS_DELETING = "deleting" NAS_STATUS_DELETE_FAILED = "delete_failed" + + NAS_UPDATE_TAGS = "update_tags" + NAS_UPDATE_TAGS_FAILED = "update_tags_fail" ) type FileSystemListInput struct { @@ -104,3 +107,8 @@ type FileSystemDetails struct { Network string Zone string } + +type FileSystemRemoteUpdateInput struct { + // 是否覆盖替换所有标签 + ReplaceTags *bool `json:"replace_tags" help:"replace all remote tags"` +} diff --git a/pkg/compute/models/filesystem.go b/pkg/compute/models/filesystem.go index 1a779ba949..49f1ff8144 100644 --- a/pkg/compute/models/filesystem.go +++ b/pkg/compute/models/filesystem.go @@ -350,7 +350,6 @@ func (self *SCloudregion) SyncFileSystems(ctx context.Context, userCred mcclient result.UpdateError(err) continue } - syncMetadata(ctx, userCred, &commondb[i], commonext[i]) localFSs = append(localFSs, commondb[i]) remoteFSs = append(remoteFSs, commonext[i]) result.Update() @@ -453,7 +452,11 @@ func (self *SFileSystem) SyncWithCloudFileSystem(ctx context.Context, userCred m } return nil }) - return errors.Wrapf(err, "db.Update") + if err != nil { + return errors.Wrapf(err, "db.Update") + } + syncMetadata(ctx, userCred, self, fs) + return nil } func (self *SCloudregion) getZoneIdBySuffix(zoneId string) (string, error) { @@ -593,3 +596,31 @@ func (manager *SFileSystemManager) DeleteExpiredPostpaids(ctx context.Context, u fss[i].StartDeleteTask(ctx, userCred, "") } } + +func (self *SFileSystem) 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 *SFileSystem) PerformRemoteUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.FileSystemRemoteUpdateInput) (jsonutils.JSONObject, error) { + return nil, self.StartRemoteUpdateTask(ctx, userCred, (input.ReplaceTags != nil && *input.ReplaceTags), "") +} + +func (self *SFileSystem) StartRemoteUpdateTask(ctx context.Context, userCred mcclient.TokenCredential, replaceTags bool, parentTaskId string) error { + data := jsonutils.NewDict() + if replaceTags { + data.Add(jsonutils.JSONTrue, "replace_tags") + } + task, err := taskman.TaskManager.NewTask(ctx, "FileSystemRemoteUpdateTask", self, userCred, data, parentTaskId, "", nil) + if err != nil { + return errors.Wrap(err, "NewTask") + } + self.SetStatus(userCred, api.NAS_UPDATE_TAGS, "StartRemoteUpdateTask") + return task.ScheduleRun(nil) +} + +func (self *SFileSystem) OnMetadataUpdated(ctx context.Context, userCred mcclient.TokenCredential) { + if len(self.ExternalId) == 0 { + return + } + self.StartRemoteUpdateTask(ctx, userCred, true, "") +} diff --git a/pkg/compute/tasks/file_system_remote_update_task.go b/pkg/compute/tasks/file_system_remote_update_task.go new file mode 100644 index 0000000000..03aac12848 --- /dev/null +++ b/pkg/compute/tasks/file_system_remote_update_task.go @@ -0,0 +1,96 @@ +// 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" + + 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/cloudprovider" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/util/logclient" +) + +type FileSystemRemoteUpdateTask struct { + taskman.STask +} + +func init() { + taskman.RegisterTask(FileSystemRemoteUpdateTask{}) +} + +func (self *FileSystemRemoteUpdateTask) taskFail(ctx context.Context, fs *models.SFileSystem, err error) { + fs.SetStatus(self.UserCred, api.NAS_UPDATE_TAGS_FAILED, err.Error()) + self.SetStageFailed(ctx, jsonutils.NewString(err.Error())) +} + +func (self *FileSystemRemoteUpdateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + fs := obj.(*models.SFileSystem) + self.SetStage("OnRemoteUpdateComplete", nil) + replaceTags := jsonutils.QueryBoolean(self.Params, "replace_tags", false) + + iFs, err := fs.GetICloudFileSystem() + if err != nil { + self.taskFail(ctx, fs, errors.Wrapf(err, "GetICloudFileSystem")) + return + } + + taskman.LocalTaskRun(self, func() (jsonutils.JSONObject, error) { + oldTags, err := iFs.GetTags() + if err != nil { + if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented { + return nil, nil + } + return nil, errors.Wrap(err, "iFs.GetTags()") + } + tags, err := fs.GetAllUserMetadata() + if err != nil { + return nil, errors.Wrapf(err, "fs.GetAllUserMetadata") + } + tagsUpdateInfo := cloudprovider.TagsUpdateInfo{OldTags: oldTags, NewTags: tags} + err = cloudprovider.SetTags(ctx, iFs, fs.ManagerId, tags, replaceTags) + if err != nil { + if errors.Cause(err) == cloudprovider.ErrNotSupported || errors.Cause(err) == cloudprovider.ErrNotImplemented { + return nil, nil + } + logclient.AddActionLogWithStartable(self, fs, logclient.ACT_UPDATE_TAGS, err, self.GetUserCred(), false) + return nil, errors.Wrap(err, "iFs.SetMetadata") + } + logclient.AddActionLogWithStartable(self, fs, logclient.ACT_UPDATE_TAGS, tagsUpdateInfo, self.GetUserCred(), true) + return nil, nil + }) +} + +func (self *FileSystemRemoteUpdateTask) OnRemoteUpdateComplete(ctx context.Context, fs *models.SFileSystem, data jsonutils.JSONObject) { + self.SetStage("OnSyncStatusComplete", nil) + models.StartResourceSyncStatusTask(ctx, self.UserCred, fs, "FileSystemSyncstatusTask", self.GetTaskId()) +} + +func (self *FileSystemRemoteUpdateTask) OnRemoteUpdateCompleteFailed(ctx context.Context, fs *models.SFileSystem, data jsonutils.JSONObject) { + self.taskFail(ctx, fs, errors.Errorf(data.String())) +} + +func (self *FileSystemRemoteUpdateTask) OnSyncStatusComplete(ctx context.Context, fs *models.SFileSystem, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} + +func (self *FileSystemRemoteUpdateTask) OnSyncStatusCompleteFailed(ctx context.Context, fs *models.SFileSystem, data jsonutils.JSONObject) { + self.SetStageFailed(ctx, data) +} diff --git a/pkg/multicloud/aliyun/aliyun.go b/pkg/multicloud/aliyun/aliyun.go index a55301c992..e20fd6bdaf 100644 --- a/pkg/multicloud/aliyun/aliyun.go +++ b/pkg/multicloud/aliyun/aliyun.go @@ -76,6 +76,7 @@ const ( ALIYUN_SERVICE_RDS = "rds" ALIYUN_SERVICE_SLB = "slb" ALIYUN_SERVICE_KVS = "kvs" + ALIYUN_SERVICE_NAS = "nas" ) var ( diff --git a/pkg/multicloud/aliyun/filesystem.go b/pkg/multicloud/aliyun/filesystem.go index 866b678ad0..aaf73370bf 100644 --- a/pkg/multicloud/aliyun/filesystem.go +++ b/pkg/multicloud/aliyun/filesystem.go @@ -356,3 +356,7 @@ func (self *SRegion) CreateFileSystem(opts *cloudprovider.FileSystemCraeteOption fsId, _ := resp.GetString("FileSystemId") return self.GetFileSystem(fsId) } + +func (self *SFileSystem) SetTags(tags map[string]string, replace bool) error { + return self.region.SetResourceTags(ALIYUN_SERVICE_NAS, "filesystem", self.FileSystemId, tags, replace) +} diff --git a/pkg/multicloud/aliyun/resource_tags.go b/pkg/multicloud/aliyun/resource_tags.go index d1c8a12bbd..d1fd247190 100644 --- a/pkg/multicloud/aliyun/resource_tags.go +++ b/pkg/multicloud/aliyun/resource_tags.go @@ -36,6 +36,8 @@ func (self *SRegion) tagRequest(serviceType, action string, params map[string]st return self.lbRequest(action, params) case ALIYUN_SERVICE_KVS: return self.kvsRequest(action, params) + case ALIYUN_SERVICE_NAS: + return self.nasRequest(action, params) default: return nil, fmt.Errorf("invalid service type") }