From 2ed4534f488fd2067642b582986e9a930885ec44 Mon Sep 17 00:00:00 2001 From: rainzm Date: Fri, 6 Nov 2020 20:55:50 +0800 Subject: [PATCH] feat: creat, delete, revert and sync operator for esxi instance snapshots --- pkg/apis/compute/snapshot_const.go | 1 + pkg/compute/models/cloudsync.go | 4 + pkg/compute/models/guest_actions.go | 7 +- pkg/compute/models/guest_queries.go | 65 +++++++++++ pkg/compute/models/instance_snapshots.go | 100 ++++++++++++++--- pkg/compute/models/regiondrivers.go | 4 + pkg/compute/regiondrivers/esxi.go | 64 +++++++++++ pkg/compute/regiondrivers/kvm.go | 104 ++++++++++++++++++ pkg/compute/regiondrivers/managedvirtual.go | 21 ++++ .../tasks/instance_snapshot_create_task.go | 69 ++++-------- .../tasks/instance_snapshot_delete_task.go | 40 +++---- .../tasks/instance_snapshot_reset_task.go | 47 ++++---- pkg/multicloud/esxi/vdisk.go | 19 +++- 13 files changed, 429 insertions(+), 116 deletions(-) diff --git a/pkg/apis/compute/snapshot_const.go b/pkg/apis/compute/snapshot_const.go index e6b4718e5f..6fc25c2b8a 100644 --- a/pkg/apis/compute/snapshot_const.go +++ b/pkg/apis/compute/snapshot_const.go @@ -45,6 +45,7 @@ const ( SNAPSHOT_POLICY_DISK_DELETE_FAILED = "delete_failed" INSTANCE_SNAPSHOT_READY = "ready" + INSTANCE_SNAPSHOT_UNKNOWN = "unknown" INSTANCE_SNAPSHOT_FAILED = "instance_snapshot_create_failed" INSTANCE_SNAPSHOT_START_DELETE = "instance_snapshot_start_delete" INSTANCE_SNAPSHOT_DELETE_FAILED = "instance_snapshot_delete_failed" diff --git a/pkg/compute/models/cloudsync.go b/pkg/compute/models/cloudsync.go index bab4c42122..9857e0a086 100644 --- a/pkg/compute/models/cloudsync.go +++ b/pkg/compute/models/cloudsync.go @@ -719,6 +719,10 @@ func syncVMPeripherals(ctx context.Context, userCred mcclient.TokenCredential, l if err != nil { log.Errorf("syncVMSecgroups error %s", err) } + result := local.SyncInstanceSnapshots(ctx, userCred, provider) + if result.IsError() { + log.Errorf("syncVMInstanceSnapshots error %v", result.AllError()) + } } func syncVMNics(ctx context.Context, userCred mcclient.TokenCredential, provider *SCloudprovider, host *SHost, localVM *SGuest, remoteVM cloudprovider.ICloudVM) error { diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 7d70da4591..7711b0e319 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -4580,6 +4580,11 @@ func (self *SGuest) AllowPerformInstanceSnapshot(ctx context.Context, return self.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, self, "instance-snapshot") } +var supportInstanceSnapshotHypervisors = []string{ + api.HYPERVISOR_KVM, + api.HYPERVISOR_ESXI, +} + func (self *SGuest) validateCreateInstanceSnapshot( ctx context.Context, userCred mcclient.TokenCredential, @@ -4587,7 +4592,7 @@ func (self *SGuest) validateCreateInstanceSnapshot( data jsonutils.JSONObject, ) (*SRegionQuota, error) { - if self.Hypervisor != api.HYPERVISOR_KVM { + if !utils.IsInStringArray(self.Hypervisor, supportInstanceSnapshotHypervisors) { return nil, httperrors.NewBadRequestError("guest hypervisor %s can't create instance snapshot", self.Hypervisor) } diff --git a/pkg/compute/models/guest_queries.go b/pkg/compute/models/guest_queries.go index ac66bc0b9a..80313c1623 100644 --- a/pkg/compute/models/guest_queries.go +++ b/pkg/compute/models/guest_queries.go @@ -21,11 +21,16 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/tristate" + "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" + "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/util/stringutils2" ) @@ -551,3 +556,63 @@ func fetchScalingGroupGuest(guestIds ...string) map[string]SScalingGroupGuest { } return ret } + +func (self *SGuest) SyncInstanceSnapshots(ctx context.Context, userCred mcclient.TokenCredential, provider *SCloudprovider) compare.SyncResult { + syncResult := compare.SyncResult{} + + extGuest, err := self.GetIVM() + if err != nil { + syncResult.Error(err) + return syncResult + } + + extSnapshots, err := extGuest.GetInstanceSnapshots() + if errors.Cause(err) == errors.ErrNotImplemented { + return syncResult + } + syncOwnerId := provider.GetOwnerId() + localSnapshots, err := self.GetInstanceSnapshots() + if err != nil { + syncResult.Error(err) + return syncResult + } + + lockman.LockClass(ctx, InstanceSnapshotManager, db.GetLockClassKey(InstanceSnapshotManager, syncOwnerId)) + defer lockman.ReleaseClass(ctx, InstanceSnapshotManager, db.GetLockClassKey(InstanceSnapshotManager, syncOwnerId)) + + removed := make([]SInstanceSnapshot, 0) + commondb := make([]SInstanceSnapshot, 0) + commonext := make([]cloudprovider.ICloudInstanceSnapshot, 0) + added := make([]cloudprovider.ICloudInstanceSnapshot, 0) + + err = compare.CompareSets(localSnapshots, extSnapshots, &removed, &commondb, &commonext, &added) + if err != nil { + syncResult.Error(err) + return syncResult + } + for i := 0; i < len(removed); i += 1 { + err = removed[i].syncRemoveCloudInstanceSnapshot(ctx, userCred) + if err != nil { + syncResult.DeleteError(err) + } else { + syncResult.Delete() + } + } + for i := 0; i < len(commondb); i += 1 { + err = commondb[i].SyncWithCloudInstanceSnapshot(ctx, userCred, commonext[i], self) + if err != nil { + syncResult.UpdateError(err) + } else { + syncResult.Update() + } + } + for i := 0; i < len(added); i += 1 { + _, err := InstanceSnapshotManager.newFromCloudInstanceSnapshot(ctx, userCred, added[i], self) + if err != nil { + syncResult.AddError(err) + } else { + syncResult.Add() + } + } + return syncResult +} diff --git a/pkg/compute/models/instance_snapshots.go b/pkg/compute/models/instance_snapshots.go index 27f1e7b7a7..a97190f75c 100644 --- a/pkg/compute/models/instance_snapshots.go +++ b/pkg/compute/models/instance_snapshots.go @@ -30,6 +30,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/util/stringutils2" @@ -49,6 +50,7 @@ func init() { type SInstanceSnapshot struct { db.SVirtualResourceBase + db.SExternalizedResourceBase // 云主机Id GuestId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" index:"true"` @@ -74,6 +76,7 @@ type SInstanceSnapshot struct { type SInstanceSnapshotManager struct { db.SVirtualResourceBaseManager + db.SExternalizedResourceBaseManager } var InstanceSnapshotManager *SInstanceSnapshotManager @@ -147,6 +150,17 @@ func (self *SInstanceSnapshot) AllowUpdateItem(ctx context.Context, userCred mcc return false } +func (self *SInstanceSnapshot) GetGuest() (*SGuest, error) { + if len(self.GuestId) == 0 { + return nil, errors.ErrNotFound + } + guest := GuestManager.FetchGuestById(self.GuestId) + if guest == nil { + return nil, errors.ErrNotFound + } + return guest, nil +} + func (self *SInstanceSnapshot) getMoreDetails(userCred mcclient.TokenCredential, out api.InstanceSnapshotDetails) api.InstanceSnapshotDetails { if guest := GuestManager.FetchGuestById(self.GuestId); guest != nil { out.Guest = guest.Name @@ -222,16 +236,11 @@ func (self *SInstanceSnapshot) StartCreateInstanceSnapshotTask( return nil } -func (manager *SInstanceSnapshotManager) CreateInstanceSnapshot( - ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, name string, autoDelete bool, -) (*SInstanceSnapshot, error) { - instanceSnapshot := &SInstanceSnapshot{} +func (manager *SInstanceSnapshotManager) fillInstanceSnapshot(userCred mcclient.TokenCredential, guest *SGuest, instanceSnapshot *SInstanceSnapshot) { instanceSnapshot.SetModelManager(manager, instanceSnapshot) - instanceSnapshot.Name = name instanceSnapshot.ProjectId = userCred.GetProjectId() instanceSnapshot.DomainId = userCred.GetProjectDomainId() instanceSnapshot.GuestId = guest.Id - instanceSnapshot.AutoDelete = autoDelete guestSchedInput := guest.ToSchedDesc() for i := 0; i < len(guestSchedInput.Disks); i++ { @@ -288,6 +297,14 @@ func (manager *SInstanceSnapshotManager) CreateInstanceSnapshot( instanceSnapshot.OsArch = guest.OsArch instanceSnapshot.ServerMetadata = serverMetadata instanceSnapshot.InstanceType = guest.InstanceType +} + +func (manager *SInstanceSnapshotManager) CreateInstanceSnapshot(ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, name string, autoDelete bool) (*SInstanceSnapshot, error) { + instanceSnapshot := &SInstanceSnapshot{} + instanceSnapshot.SetModelManager(manager, instanceSnapshot) + instanceSnapshot.Name = name + instanceSnapshot.AutoDelete = autoDelete + manager.fillInstanceSnapshot(userCred, guest, instanceSnapshot) err := manager.TableSpec().Insert(ctx, instanceSnapshot) if err != nil { return nil, err @@ -303,13 +320,15 @@ func (self *SInstanceSnapshot) ToInstanceCreateInput( return nil, errors.Wrap(err, "unmarshal sched input") } - isjs := make([]SInstanceSnapshotJoint, 0) - err := InstanceSnapshotJointManager.Query().Equals("instance_snapshot_id", self.Id).Asc("disk_index").All(&isjs) - if err != nil { - return nil, errors.Wrap(err, "fetch instance snapshots") - } - for i := 0; i < len(serverConfig.Disks); i++ { - serverConfig.Disks[i].SnapshotId = isjs[serverConfig.Disks[i].Index].SnapshotId + if len(self.ExternalId) > 0 { + isjs := make([]SInstanceSnapshotJoint, 0) + err := InstanceSnapshotJointManager.Query().Equals("instance_snapshot_id", self.Id).Asc("disk_index").All(&isjs) + if err != nil { + return nil, errors.Wrap(err, "fetch instance snapshots") + } + for i := 0; i < len(serverConfig.Disks); i++ { + serverConfig.Disks[i].SnapshotId = isjs[serverConfig.Disks[i].Index].SnapshotId + } } sourceInput.Disks = serverConfig.Disks if sourceInput.VmemSize == 0 { @@ -419,3 +438,58 @@ func (self *SInstanceSnapshot) DecRefCount(ctx context.Context, userCred mcclien } return err } + +func (is *SInstanceSnapshot) syncRemoveCloudInstanceSnapshot(ctx context.Context, userCred mcclient.TokenCredential) error { + lockman.LockObject(ctx, is) + defer lockman.ReleaseObject(ctx, is) + + err := is.ValidateDeleteCondition(ctx) + if err != nil { + err = is.SetStatus(userCred, api.INSTANCE_SNAPSHOT_UNKNOWN, "sync to delete") + } else { + err = is.RealDelete(ctx, userCred) + } + return err +} + +func (is *SInstanceSnapshot) SyncWithCloudInstanceSnapshot(ctx context.Context, userCred mcclient.TokenCredential, ext cloudprovider.ICloudInstanceSnapshot, guest *SGuest) error { + diff, err := db.UpdateWithLock(ctx, is, func() error { + is.Status = ext.GetStatus() + InstanceSnapshotManager.fillInstanceSnapshot(userCred, guest, is) + return nil + }) + if err != nil { + return err + } + db.OpsLog.LogSyncUpdate(is, diff, userCred) + return nil +} + +func (manager *SInstanceSnapshotManager) newFromCloudInstanceSnapshot(ctx context.Context, userCred mcclient.TokenCredential, extSnapshot cloudprovider.ICloudInstanceSnapshot, guest *SGuest) (*SInstanceSnapshot, error) { + instanceSnapshot := SInstanceSnapshot{} + instanceSnapshot.SetModelManager(manager, &instanceSnapshot) + newName, err := db.GenerateName(manager, nil, extSnapshot.GetName()) + if err == nil { + instanceSnapshot.Name = extSnapshot.GetName() + } else { + instanceSnapshot.Name = newName + } + instanceSnapshot.ExternalId = extSnapshot.GetGlobalId() + instanceSnapshot.Status = extSnapshot.GetStatus() + manager.fillInstanceSnapshot(userCred, guest, &instanceSnapshot) + err = manager.TableSpec().Insert(ctx, &instanceSnapshot) + if err != nil { + return nil, err + } + db.OpsLog.LogEvent(&instanceSnapshot, db.ACT_CREATE, instanceSnapshot.GetShortDesc(ctx), userCred) + return &instanceSnapshot, nil +} + +func (self *SInstanceSnapshot) GetRegionDriver() IRegionDriver { + guest, _ := self.GetGuest() + var provider string + if guest != nil { + provider = guest.GetHost().GetProviderName() + } + return GetRegionDriver(provider) +} diff --git a/pkg/compute/models/regiondrivers.go b/pkg/compute/models/regiondrivers.go index 2ebfc05914..c44d4545e3 100644 --- a/pkg/compute/models/regiondrivers.go +++ b/pkg/compute/models/regiondrivers.go @@ -115,6 +115,10 @@ type IRegionDriver interface { OnDiskReset(ctx context.Context, userCred mcclient.TokenCredential, disk *SDisk, snapshot *SSnapshot, data jsonutils.JSONObject) error OnSnapshotDelete(ctx context.Context, snapshot *SSnapshot, task taskman.ITask, data jsonutils.JSONObject) error + RequestCreateInstanceSnapshot(ctx context.Context, guest *SGuest, isp *SInstanceSnapshot, task taskman.ITask, params *jsonutils.JSONDict) error + RequestDeleteInstanceSnapshot(ctx context.Context, isp *SInstanceSnapshot, task taskman.ITask) error + RequestResetToInstanceSnapshot(ctx context.Context, guest *SGuest, isp *SInstanceSnapshot, task taskman.ITask, params *jsonutils.JSONDict) error + //Nat gateway DealNatGatewaySpec(spec string) string RequestBindIPToNatgateway(ctx context.Context, task taskman.ITask, natgateway *SNatGateway, eipID string) error diff --git a/pkg/compute/regiondrivers/esxi.go b/pkg/compute/regiondrivers/esxi.go index cbba50ed95..9aa3b321d4 100644 --- a/pkg/compute/regiondrivers/esxi.go +++ b/pkg/compute/regiondrivers/esxi.go @@ -19,8 +19,11 @@ import ( "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" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" @@ -54,3 +57,64 @@ func (self *SEsxiRegionDriver) ValidateCreateLoadbalancerCertificateData(ctx con func (self *SEsxiRegionDriver) ValidateCreateSnapshotData(ctx context.Context, userCred mcclient.TokenCredential, disk *models.SDisk, storage *models.SStorage, input *api.SnapshotCreateInput) error { return fmt.Errorf("%s does not support creating snapshot", self.GetProvider()) } + +func (self *SEsxiRegionDriver) RequestCreateInstanceSnapshot(ctx context.Context, guest *models.SGuest, isp *models.SInstanceSnapshot, task taskman.ITask, params *jsonutils.JSONDict) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + ivm, err := guest.GetIVM() + if err != nil { + return nil, errors.Wrap(err, "unable to GetIVM") + } + cloudSP, err := ivm.CreateInstanceSnapshot(ctx, isp.GetName(), isp.Description) + if err != nil { + return nil, errors.Wrap(err, "unable to CreateInstanceSnapshot") + } + _, err = db.Update(isp, func() error { + isp.SetExternalId(cloudSP.GetGlobalId()) + return nil + }) + return nil, err + }) + return nil +} + +func (self *SEsxiRegionDriver) RequestDeleteInstanceSnapshot(ctx context.Context, isp *models.SInstanceSnapshot, task taskman.ITask) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + guest, err := isp.GetGuest() + if err != nil { + return nil, errors.Wrap(err, "GetGuest") + } + ivm, err := guest.GetIVM() + if err != nil { + return nil, errors.Wrap(err, "unable to GetIVM") + } + id := isp.GetExternalId() + if len(id) == 0 { + return nil, nil + } + cloudSP, err := ivm.GetInstanceSnapshot(id) + if err != nil { + return nil, errors.Wrap(err, "unable to GetInstanceSnapshot") + } + err = cloudSP.Delete() + if err != nil { + return nil, errors.Wrap(err, "unable to delete cloud instance snapshot") + } + return nil, nil + }) + return nil +} + +func (self *SEsxiRegionDriver) RequestResetToInstanceSnapshot(ctx context.Context, guest *models.SGuest, isp *models.SInstanceSnapshot, task taskman.ITask, params *jsonutils.JSONDict) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + ivm, err := guest.GetIVM() + if err != nil { + return nil, errors.Wrap(err, "unable to GetIVM") + } + err = ivm.ResetToInstanceSnapshot(ctx, isp.GetExternalId()) + if err != nil { + return nil, errors.Wrap(err, "unable to ResetToInstanceSnapshot") + } + return nil, nil + }) + return nil +} diff --git a/pkg/compute/regiondrivers/kvm.go b/pkg/compute/regiondrivers/kvm.go index bb66bc301a..98bd2cc3b4 100644 --- a/pkg/compute/regiondrivers/kvm.go +++ b/pkg/compute/regiondrivers/kvm.go @@ -19,6 +19,7 @@ import ( "database/sql" "fmt" "regexp" + "strconv" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -27,6 +28,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/cloudcommon/validators" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -965,6 +967,61 @@ func (self *SKVMRegionDriver) RequestDeleteSnapshot(ctx context.Context, snapsho return models.GetStorageDriver(storage.StorageType).RequestDeleteSnapshot(ctx, snapshot, task) } +func (self *SKVMRegionDriver) RequestDeleteInstanceSnapshot(ctx context.Context, isp *models.SInstanceSnapshot, task taskman.ITask) error { + snapshots, err := isp.GetSnapshots() + if err != nil { + return err + } + if len(snapshots) == 0 { + task.SetStage("OnInstanceSnapshotDelete", nil) + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + return nil, nil + }) + return nil + } + + params := jsonutils.NewDict() + params.Set("del_snapshot_id", jsonutils.NewString(snapshots[0].Id)) + task.SetStage("OnKvmSnapshotDelete", params) + err = snapshots[0].StartSnapshotDeleteTask(ctx, task.GetUserCred(), false, task.GetTaskId()) + if err != nil { + return err + } + return nil +} + +func (self *SKVMRegionDriver) RequestResetToInstanceSnapshot(ctx context.Context, guest *models.SGuest, isp *models.SInstanceSnapshot, task taskman.ITask, params *jsonutils.JSONDict) error { + disks := guest.GetDisks() + diskIndexI64, err := params.Int("disk_index") + if err != nil { + return errors.Wrap(err, "get 'disk_index' from params") + } + diskIndex := int(diskIndexI64) + if diskIndex >= len(disks) { + task.SetStage("OnInstanceSnapshotReset", nil) + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + return nil, nil + }) + return nil + } + + isj, err := isp.GetInstanceSnapshotJointAt(diskIndex) + if err != nil { + return err + } + + params = jsonutils.NewDict() + params.Set("disk_index", jsonutils.NewInt(int64(diskIndex))) + task.SetStage("OnKvmDiskReset", params) + + disk := disks[diskIndex].GetDisk() + err = disk.StartResetDisk(ctx, task.GetUserCred(), isj.SnapshotId, false, guest, task.GetTaskId()) + if err != nil { + return err + } + return nil +} + func (self *SKVMRegionDriver) ValidateCreateSnapshotData(ctx context.Context, userCred mcclient.TokenCredential, disk *models.SDisk, storage *models.SStorage, input *api.SnapshotCreateInput) error { host := storage.GetMasterHost() if host == nil { @@ -981,6 +1038,53 @@ func (self *SKVMRegionDriver) RequestCreateSnapshot(ctx context.Context, snapsho return models.GetStorageDriver(storage.StorageType).RequestCreateSnapshot(ctx, snapshot, task) } +func (self *SKVMRegionDriver) RequestCreateInstanceSnapshot(ctx context.Context, guest *models.SGuest, isp *models.SInstanceSnapshot, task taskman.ITask, params *jsonutils.JSONDict) error { + disks := guest.GetDisks() + diskIndexI64, err := params.Int("disk_index") + if err != nil { + return errors.Wrap(err, "get 'disk_index' from params") + } + diskIndex := int(diskIndexI64) + if diskIndex >= len(disks) { + task.SetStage("OnInstanceSnapshot", nil) + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + return nil, nil + }) + return nil + } + + lockman.LockClass(ctx, models.SnapshotManager, task.GetUserCred().GetProjectId()) + defer lockman.ReleaseClass(ctx, models.SnapshotManager, task.GetUserCred().GetProjectId()) + + snapshotName, err := db.GenerateName(models.SnapshotManager, task.GetUserCred(), + fmt.Sprintf("%s-%s", isp.Name, rand.String(8))) + if err != nil { + return errors.Wrap(err, "Generate snapshot name") + } + + snapshot, err := models.SnapshotManager.CreateSnapshot( + ctx, task.GetUserCred(), api.SNAPSHOT_MANUAL, disks[diskIndex].DiskId, + guest.Id, "", snapshotName, -1) + if err != nil { + return err + } + + err = models.InstanceSnapshotJointManager.CreateJoint(ctx, isp.Id, snapshot.Id, int8(diskIndex)) + if err != nil { + return err + } + + params = jsonutils.NewDict() + params.Set("disk_index", jsonutils.NewInt(int64(diskIndex))) + params.Set(strconv.Itoa(diskIndex), jsonutils.NewString(snapshot.Id)) + task.SetStage("OnKvmDiskSnapshot", params) + + if err := snapshot.StartSnapshotCreateTask(ctx, task.GetUserCred(), nil, task.GetTaskId()); err != nil { + return err + } + return nil +} + func (self *SKVMRegionDriver) SnapshotIsOutOfChain(disk *models.SDisk) bool { storage := disk.GetStorage() return models.GetStorageDriver(storage.StorageType).SnapshotIsOutOfChain(disk) diff --git a/pkg/compute/regiondrivers/managedvirtual.go b/pkg/compute/regiondrivers/managedvirtual.go index 7041460241..5425398a14 100644 --- a/pkg/compute/regiondrivers/managedvirtual.go +++ b/pkg/compute/regiondrivers/managedvirtual.go @@ -1341,6 +1341,27 @@ func (self *SManagedVirtualizationRegionDriver) RequestCreateSnapshot(ctx contex return nil } +func (self *SManagedVirtualizationRegionDriver) RequestCreateInstanceSnapshot(ctx context.Context, guest *models.SGuest, isp *models.SInstanceSnapshot, task taskman.ITask, params *jsonutils.JSONDict) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + return nil, nil + }) + return nil +} + +func (self *SManagedVirtualizationRegionDriver) RequestDeleteInstanceSnapshot(ctx context.Context, isp *models.SInstanceSnapshot, task taskman.ITask) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + return nil, nil + }) + return nil +} + +func (self *SManagedVirtualizationRegionDriver) RequestResetToInstanceSnapshot(ctx context.Context, guest *models.SGuest, isp *models.SInstanceSnapshot, task taskman.ITask, params *jsonutils.JSONDict) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + return nil, nil + }) + return nil +} + func (self *SManagedVirtualizationRegionDriver) GetDiskResetParams(snapshot *models.SSnapshot) *jsonutils.JSONDict { params := jsonutils.NewDict() params.Set("snapshot_id", jsonutils.NewString(snapshot.ExternalId)) diff --git a/pkg/compute/tasks/instance_snapshot_create_task.go b/pkg/compute/tasks/instance_snapshot_create_task.go index 5bbc733de8..3f2a438855 100644 --- a/pkg/compute/tasks/instance_snapshot_create_task.go +++ b/pkg/compute/tasks/instance_snapshot_create_task.go @@ -16,20 +16,16 @@ package tasks import ( "context" - "fmt" - "strconv" "yunion.io/x/jsonutils" "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" - "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" - "yunion.io/x/onecloud/pkg/util/rand" ) type InstanceSnapshotCreateTask struct { @@ -88,55 +84,16 @@ func (self *InstanceSnapshotCreateTask) OnInit( isp := obj.(*models.SInstanceSnapshot) guest := models.GuestManager.FetchGuestById(isp.GuestId) - - self.GuestDiskCreateSnapshot(ctx, isp, guest, 0) -} - -func (self *InstanceSnapshotCreateTask) GuestDiskCreateSnapshot( - ctx context.Context, isp *models.SInstanceSnapshot, guest *models.SGuest, diskIndex int) { - - disks := guest.GetDisks() - if diskIndex >= len(disks) { - self.taskComplete(ctx, isp, guest, nil) - return - } - - lockman.LockClass(ctx, models.SnapshotManager, self.UserCred.GetProjectId()) - defer lockman.ReleaseClass(ctx, models.SnapshotManager, self.UserCred.GetProjectId()) - - snapshotName, err := db.GenerateName(models.SnapshotManager, self.UserCred, - fmt.Sprintf("%s-%s", isp.Name, rand.String(8))) - if err != nil { - self.taskFail(ctx, isp, guest, jsonutils.NewString(fmt.Sprintf("Generate snapshot name %s", err))) - return - } - - snapshot, err := models.SnapshotManager.CreateSnapshot( - ctx, self.UserCred, compute.SNAPSHOT_MANUAL, disks[diskIndex].DiskId, - guest.Id, "", snapshotName, -1) - if err != nil { - self.taskFail(ctx, isp, guest, jsonutils.NewString(err.Error())) - return - } - - err = models.InstanceSnapshotJointManager.CreateJoint(ctx, isp.Id, snapshot.Id, int8(diskIndex)) - if err != nil { - self.taskFail(ctx, isp, guest, jsonutils.NewString(err.Error())) - return - } - + self.SetStage("OnInstanceSnapshot", nil) params := jsonutils.NewDict() - params.Set("disk_index", jsonutils.NewInt(int64(diskIndex))) - params.Set(strconv.Itoa(diskIndex), jsonutils.NewString(snapshot.Id)) - self.SetStage("OnDiskSnapshot", params) - - if err := snapshot.StartSnapshotCreateTask(ctx, self.UserCred, nil, self.Id); err != nil { + params.Set("disk_index", jsonutils.NewInt(0)) + if err := isp.GetRegionDriver().RequestCreateInstanceSnapshot(ctx, guest, isp, self, params); err != nil { self.taskFail(ctx, isp, guest, jsonutils.NewString(err.Error())) return } } -func (self *InstanceSnapshotCreateTask) OnDiskSnapshot( +func (self *InstanceSnapshotCreateTask) OnKvmDiskSnapshot( ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { guest := models.GuestManager.FetchGuestById(isp.GuestId) @@ -147,10 +104,24 @@ func (self *InstanceSnapshotCreateTask) OnDiskSnapshot( return } - self.GuestDiskCreateSnapshot(ctx, isp, guest, int(diskIndex+1)) + params := jsonutils.NewDict() + params.Set("disk_index", jsonutils.NewInt(diskIndex+1)) + if err := isp.GetRegionDriver().RequestCreateInstanceSnapshot(ctx, guest, isp, self, params); err != nil { + self.taskFail(ctx, isp, guest, jsonutils.NewString(err.Error())) + return + } } -func (self *InstanceSnapshotCreateTask) OnDiskSnapshotFailed( +func (self *InstanceSnapshotCreateTask) OnKvmDiskSnapshotFailed( ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { self.taskFail(ctx, isp, nil, data) } + +func (self *InstanceSnapshotCreateTask) OnInstanceSnapshot(ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { + guest, _ := isp.GetGuest() + self.taskComplete(ctx, isp, guest, data) +} + +func (self *InstanceSnapshotCreateTask) OnInstanceSnapshotFailed(ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { + self.taskFail(ctx, isp, nil, data) +} diff --git a/pkg/compute/tasks/instance_snapshot_delete_task.go b/pkg/compute/tasks/instance_snapshot_delete_task.go index 782f43e89f..79a35cf30a 100644 --- a/pkg/compute/tasks/instance_snapshot_delete_task.go +++ b/pkg/compute/tasks/instance_snapshot_delete_task.go @@ -55,33 +55,14 @@ func (self *InstanceSnapshotDeleteTask) OnInit( ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { isp := obj.(*models.SInstanceSnapshot) - self.StartSnapshotDelete(ctx, isp) -} - -func (self *InstanceSnapshotDeleteTask) StartSnapshotDelete( - ctx context.Context, isp *models.SInstanceSnapshot) { - - snapshots, err := isp.GetSnapshots() - if err != nil { - self.taskFail(ctx, isp, jsonutils.NewString(err.Error())) - return - } - if len(snapshots) == 0 { - self.taskComplete(ctx, isp, nil) - return - } - - params := jsonutils.NewDict() - params.Set("del_snapshot_id", jsonutils.NewString(snapshots[0].Id)) - self.SetStage("OnSnapshotDelete", params) - err = snapshots[0].StartSnapshotDeleteTask(ctx, self.UserCred, false, self.Id) - if err != nil { + self.SetStage("OnInstanceSnapshotDelete", nil) + if err := isp.GetRegionDriver().RequestDeleteInstanceSnapshot(ctx, isp, self); err != nil { self.taskFail(ctx, isp, jsonutils.NewString(err.Error())) return } } -func (self *InstanceSnapshotDeleteTask) OnSnapshotDelete( +func (self *InstanceSnapshotDeleteTask) OnKvmSnapshotDelete( ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { snapshotId, _ := self.Params.GetString("del_snapshot_id") // detach snapshot and instance @@ -98,11 +79,22 @@ func (self *InstanceSnapshotDeleteTask) OnSnapshotDelete( self.taskFail(ctx, isp, jsonutils.NewString(err.Error())) return } - self.StartSnapshotDelete(ctx, isp) + if err := isp.GetRegionDriver().RequestDeleteInstanceSnapshot(ctx, isp, self); err != nil { + self.taskFail(ctx, isp, jsonutils.NewString(err.Error())) + return + } } -func (self *InstanceSnapshotDeleteTask) OnSnapshotDeleteFailed( +func (self *InstanceSnapshotDeleteTask) OnKvmSnapshotDeleteFailed( ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { + self.taskFail(ctx, isp, data) +} +func (self *InstanceSnapshotDeleteTask) OnInstanceSnapshotDelete(ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { + self.taskComplete(ctx, isp, data) + +} + +func (self *InstanceSnapshotDeleteTask) OnInstanceSnapshotDeleteFailed(ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { self.taskFail(ctx, isp, data) } diff --git a/pkg/compute/tasks/instance_snapshot_reset_task.go b/pkg/compute/tasks/instance_snapshot_reset_task.go index 63943006d6..83d5dcd480 100644 --- a/pkg/compute/tasks/instance_snapshot_reset_task.go +++ b/pkg/compute/tasks/instance_snapshot_reset_task.go @@ -68,37 +68,16 @@ func (self *InstanceSnapshotResetTask) OnInit( isp := obj.(*models.SInstanceSnapshot) guest := models.GuestManager.FetchGuestById(isp.GuestId) - self.GuestDiskResetTask(ctx, isp, guest, 0) -} - -func (self *InstanceSnapshotResetTask) GuestDiskResetTask( - ctx context.Context, isp *models.SInstanceSnapshot, guest *models.SGuest, diskIndex int) { - - disks := guest.GetDisks() - if diskIndex >= len(disks) { - self.taskComplete(ctx, isp, guest, nil) - return - } - - isj, err := isp.GetInstanceSnapshotJointAt(diskIndex) - if err != nil { - self.taskFail(ctx, isp, guest, jsonutils.NewString(err.Error())) - return - } - + self.SetStage("OnInstanceSnapshotReset", nil) params := jsonutils.NewDict() - params.Set("disk_index", jsonutils.NewInt(int64(diskIndex))) - self.SetStage("OnDiskReset", params) - - disk := disks[diskIndex].GetDisk() - err = disk.StartResetDisk(ctx, self.UserCred, isj.SnapshotId, false, guest, self.Id) - if err != nil { + params.Set("disk_index", jsonutils.NewInt(0)) + if err := isp.GetRegionDriver().RequestResetToInstanceSnapshot(ctx, guest, isp, self, params); err != nil { self.taskFail(ctx, isp, guest, jsonutils.NewString(err.Error())) return } } -func (self *InstanceSnapshotResetTask) OnDiskReset( +func (self *InstanceSnapshotResetTask) OnKvmDiskReset( ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { guest := models.GuestManager.FetchGuestById(isp.GuestId) @@ -108,10 +87,24 @@ func (self *InstanceSnapshotResetTask) OnDiskReset( self.taskFail(ctx, isp, guest, jsonutils.NewString(err.Error())) return } - self.GuestDiskResetTask(ctx, isp, guest, int(diskIndex+1)) + params := jsonutils.NewDict() + params.Set("disk_index", jsonutils.NewInt(diskIndex+1)) + if err := isp.GetRegionDriver().RequestResetToInstanceSnapshot(ctx, guest, isp, self, params); err != nil { + self.taskFail(ctx, isp, guest, jsonutils.NewString(err.Error())) + return + } } -func (self *InstanceSnapshotResetTask) OnDiskResetFailed( +func (self *InstanceSnapshotResetTask) OnKvmDiskResetFailed( ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { self.taskFail(ctx, isp, nil, data) } + +func (self *InstanceSnapshotResetTask) OnInstanceSnapshotReset(ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { + guest, _ := isp.GetGuest() + self.taskComplete(ctx, isp, guest, data) +} + +func (self *InstanceSnapshotResetTask) OnInstanceSnapshotResetFailed(ctx context.Context, isp *models.SInstanceSnapshot, data jsonutils.JSONObject) { + self.taskFail(ctx, isp, nil, data) +} diff --git a/pkg/multicloud/esxi/vdisk.go b/pkg/multicloud/esxi/vdisk.go index 0a89b8f856..ac15c9f793 100644 --- a/pkg/multicloud/esxi/vdisk.go +++ b/pkg/multicloud/esxi/vdisk.go @@ -288,6 +288,22 @@ func (disk *SVirtualDisk) GetId() string { return backing.GetUuid() } +func (disk *SVirtualDisk) MatchId(id string) bool { + vmid := disk.vm.GetGlobalId() + if !strings.HasPrefix(id, vmid) { + return false + } + backingUuid := id[len(vmid)+1:] + backing := disk.getBackingInfo() + for backing != nil { + if backing.GetUuid() == backingUuid { + return true + } + backing = backing.GetParent() + } + return false +} + func (disk *SVirtualDisk) GetName() string { backing := disk.getBackingInfo() return path.Base(backing.GetFileName()) @@ -362,8 +378,7 @@ func (disk *SVirtualDisk) GetTemplateId() string { } func (disk *SVirtualDisk) GetDiskType() string { - backing := disk.getBackingInfo() - if backing.GetParent() != nil { + if disk.index == 0 { return api.DISK_TYPE_SYS } return api.DISK_TYPE_DATA