diff --git a/pkg/compute/guestdrivers/esxi.go b/pkg/compute/guestdrivers/esxi.go index 316da7cb05..9ba7967f21 100644 --- a/pkg/compute/guestdrivers/esxi.go +++ b/pkg/compute/guestdrivers/esxi.go @@ -30,7 +30,6 @@ import ( "yunion.io/x/pkg/util/httputils" "yunion.io/x/pkg/util/rbacscope" "yunion.io/x/pkg/utils" - "yunion.io/x/sqlchemy" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" @@ -601,80 +600,6 @@ func (self *SESXiGuestDriver) IsSupportLiveMigrate() bool { return true } -func (self *SESXiGuestDriver) RequestMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential, input api.GuestMigrateInput, task taskman.ITask) error { - return self.RequestLiveMigrate(ctx, guest, userCred, api.GuestLiveMigrateInput{PreferHostId: input.PreferHostId}, task) -} - -func (self *SESXiGuestDriver) RequestLiveMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential, input api.GuestLiveMigrateInput, task taskman.ITask) error { - taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { - iVM, err := guest.GetIVM(ctx) - if err != nil { - return nil, errors.Wrap(err, "guest.GetIVM") - } - iHost, err := models.HostManager.FetchById(input.PreferHostId) - if err != nil { - return nil, errors.Wrapf(err, "models.HostManager.FetchById(%s)", input.PreferHostId) - } - host := iHost.(*models.SHost) - hostExternalId := host.ExternalId - if err = iVM.LiveMigrateVM(hostExternalId); err != nil { - return nil, errors.Wrapf(err, "iVM.LiveMigrateVM(%s)", hostExternalId) - } - hostExternalId = iVM.GetIHostId() - if hostExternalId == "" { - return nil, errors.Wrap(fmt.Errorf("empty hostExternalId"), "iVM.GetIHostId()") - } - iHost, err = db.FetchByExternalIdAndManagerId(models.HostManager, hostExternalId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery { - if host, _ := guest.GetHost(); host != nil { - return q.Equals("manager_id", host.ManagerId) - } - return q - }) - if err != nil { - return nil, errors.Wrapf(err, "db.FetchByExternalId(models.HostManager,%s)", hostExternalId) - } - host = iHost.(*models.SHost) - _, err = db.Update(guest, func() error { - guest.HostId = host.GetId() - return nil - }) - if err != nil { - return nil, errors.Wrap(err, "db.Update guest.hostId") - } - disks, err := guest.GetDisks() - if err != nil { - return nil, errors.Wrapf(err, "GetDisks") - } - iRegion, err := host.GetIRegion(ctx) - if err != nil { - return nil, errors.Wrapf(err, "GetIRegion") - } - for i := range disks { - iDisk, err := iRegion.GetIDiskById(disks[i].ExternalId) - if err != nil { - return nil, errors.Wrapf(err, "GetIDisk(%s)", disks[i].ExternalId) - } - iStorage, err := db.FetchByExternalIdAndManagerId(models.StorageManager, iDisk.GetIStorageId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery { - hcs := models.HoststorageManager.Query().SubQuery() - return q.Join(hcs, sqlchemy.Equals(hcs.Field("storage_id"), q.Field("id"))).Filter(sqlchemy.Equals(hcs.Field("host_id"), host.GetId())) - }) - if err != nil { - return nil, errors.Wrapf(err, "FetchStorageByExternalId(%s)", iDisk.GetIStorageId()) - } - storage := iStorage.(*models.SStorage) - _, err = db.Update(&disks[i], func() error { - disks[i].StorageId = storage.Id - return nil - }) - if err != nil { - return nil, errors.Wrapf(err, "db.Update disk %s storageid", disks[i].Name) - } - } - return nil, nil - }) - return nil -} - func (self *SESXiGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error { ivm, err := guest.GetIVM(ctx) if err != nil { diff --git a/pkg/compute/guestdrivers/managedvirtual.go b/pkg/compute/guestdrivers/managedvirtual.go index f0fdebf728..94c6082ad6 100644 --- a/pkg/compute/guestdrivers/managedvirtual.go +++ b/pkg/compute/guestdrivers/managedvirtual.go @@ -18,6 +18,7 @@ import ( "context" "fmt" "math" + "strings" "sync" "time" @@ -1324,6 +1325,7 @@ func GetCloudVMStatus(vm cloudprovider.ICloudVM) string { status = cloudprovider.CloudVMStatusDeploying case api.VM_SUSPEND: status = cloudprovider.CloudVMStatusSuspend + case api.VM_MIGRATING, api.VM_START_MIGRATE: default: status = cloudprovider.CloudVMStatusOther } @@ -1331,6 +1333,84 @@ func GetCloudVMStatus(vm cloudprovider.ICloudVM) string { return status } +func (self *SManagedVirtualizedGuestDriver) RequestMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential, input api.GuestMigrateInput, task taskman.ITask) error { + return self.requestMigrate(ctx, guest, userCred, api.GuestLiveMigrateInput{PreferHostId: input.PreferHostId}, task, false) +} + +func (self *SManagedVirtualizedGuestDriver) RequestLiveMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential, input api.GuestLiveMigrateInput, task taskman.ITask) error { + return self.requestMigrate(ctx, guest, userCred, input, task, true) +} + +func (self *SManagedVirtualizedGuestDriver) requestMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential, input api.GuestLiveMigrateInput, task taskman.ITask, isLive bool) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + iVM, err := guest.GetIVM(ctx) + if err != nil { + return nil, errors.Wrap(err, "guest.GetIVM") + } + iHost, err := models.HostManager.FetchById(input.PreferHostId) + if err != nil { + return nil, errors.Wrapf(err, "FetchById(%s)", input.PreferHostId) + } + host := iHost.(*models.SHost) + hostExternalId := host.ExternalId + if isLive { + err = iVM.LiveMigrateVM(hostExternalId) + } else { + err = iVM.MigrateVM(hostExternalId) + } + if err != nil { + return nil, errors.Wrapf(err, "Migrate (%s)", hostExternalId) + } + err = cloudprovider.Wait(time.Second*10, time.Hour*1, func() (bool, error) { + err = iVM.Refresh() + if err != nil { + return false, err + } + vmStatus := iVM.GetStatus() + if vmStatus == api.VM_UNKNOWN || strings.Contains(vmStatus, "fail") { + return false, errors.Wrapf(cloudprovider.ErrInvalidStatus, vmStatus) + } + if !utils.IsInStringArray(vmStatus, []string{api.VM_RUNNING, api.VM_READY}) { + return false, nil + } + hostId := iVM.GetIHostId() + if len(hostId) > 0 && hostId != hostExternalId { + hostExternalId = hostId + return true, nil + } + return false, nil + }) + if err != nil { + return nil, errors.Wrapf(err, "wait host change") + } + iHost, err = db.FetchByExternalIdAndManagerId(models.HostManager, hostExternalId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery { + if host, _ := guest.GetHost(); host != nil { + return q.Equals("manager_id", host.ManagerId) + } + return q + }) + if err != nil { + return nil, errors.Wrapf(err, "fetch host %s", hostExternalId) + } + host = iHost.(*models.SHost) + _, err = db.Update(guest, func() error { + guest.HostId = host.GetId() + return nil + }) + if err != nil { + return nil, errors.Wrap(err, "update hostId") + } + provider := host.GetCloudprovider() + driver, err := provider.GetProvider(ctx) + if err != nil { + return nil, err + } + models.SyncVMPeripherals(ctx, userCred, guest, iVM, host, provider, driver) + return nil, nil + }) + return nil +} + func (drv *SManagedVirtualizedGuestDriver) RequestConvertPublicipToEip(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, task taskman.ITask) error { taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { iVM, err := guest.GetIVM(ctx) diff --git a/pkg/compute/guestdrivers/openstack.go b/pkg/compute/guestdrivers/openstack.go index c5121503fc..41459061d4 100644 --- a/pkg/compute/guestdrivers/openstack.go +++ b/pkg/compute/guestdrivers/openstack.go @@ -26,7 +26,6 @@ import ( "yunion.io/x/pkg/util/billing" "yunion.io/x/pkg/util/rbacscope" "yunion.io/x/pkg/utils" - "yunion.io/x/sqlchemy" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" @@ -387,91 +386,3 @@ func (self *SOpenStackGuestDriver) IsSupportMigrate() bool { func (self *SOpenStackGuestDriver) IsSupportLiveMigrate() bool { return true } - -func (self *SOpenStackGuestDriver) RequestMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential, input api.GuestMigrateInput, task taskman.ITask) error { - taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { - iVM, err := guest.GetIVM(ctx) - if err != nil { - return nil, errors.Wrap(err, "guest.GetIVM") - } - hostExternalId := "" - if input.PreferHostId != "" { - iHost, err := models.HostManager.FetchById(input.PreferHostId) - if err != nil { - return nil, errors.Wrapf(err, "models.HostManager.FetchById(%s)", input.PreferHostId) - } - host := iHost.(*models.SHost) - hostExternalId = host.ExternalId - } - if err = iVM.MigrateVM(hostExternalId); err != nil { - return nil, errors.Wrapf(err, "iVM.MigrateVM(%s)", hostExternalId) - } - hostExternalId = iVM.GetIHostId() - if hostExternalId == "" { - return nil, errors.Wrap(fmt.Errorf("empty hostExternalId"), "iVM.GetIHostId()") - } - iHost, err := db.FetchByExternalIdAndManagerId(models.HostManager, hostExternalId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery { - if host, _ := guest.GetHost(); host != nil { - return q.Equals("manager_id", host.ManagerId) - } - return q - }) - if err != nil { - return nil, errors.Wrapf(err, "db.FetchByExternalId(models.HostManager,%s)", hostExternalId) - } - host := iHost.(*models.SHost) - _, err = db.Update(guest, func() error { - guest.HostId = host.GetId() - return nil - }) - if err != nil { - return nil, errors.Wrap(err, "db.Update guest.hostId") - } - return nil, nil - }) - return nil -} - -func (self *SOpenStackGuestDriver) RequestLiveMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential, input api.GuestLiveMigrateInput, task taskman.ITask) error { - taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { - iVM, err := guest.GetIVM(ctx) - if err != nil { - return nil, errors.Wrap(err, "guest.GetIVM") - } - hostExternalId := "" - if input.PreferHostId != "" { - iHost, err := models.HostManager.FetchById(input.PreferHostId) - if err != nil { - return nil, errors.Wrapf(err, "models.HostManager.FetchById(%s)", input.PreferHostId) - } - host := iHost.(*models.SHost) - hostExternalId = host.ExternalId - } - if err = iVM.LiveMigrateVM(hostExternalId); err != nil { - return nil, errors.Wrapf(err, "iVM.LiveMigrateVM(%s)", hostExternalId) - } - hostExternalId = iVM.GetIHostId() - if hostExternalId == "" { - return nil, errors.Wrap(fmt.Errorf("empty hostExternalId"), "iVM.GetIHostId()") - } - iHost, err := db.FetchByExternalIdAndManagerId(models.HostManager, hostExternalId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery { - if host, _ := guest.GetHost(); host != nil { - return q.Equals("manager_id", host.ManagerId) - } - return q - }) - if err != nil { - return nil, errors.Wrapf(err, "db.FetchByExternalId(models.HostManager,%s)", hostExternalId) - } - host := iHost.(*models.SHost) - _, err = db.Update(guest, func() error { - guest.HostId = host.GetId() - return nil - }) - if err != nil { - return nil, errors.Wrap(err, "db.Update guest.hostId") - } - return nil, nil - }) - return nil -} diff --git a/pkg/compute/models/cloudsync.go b/pkg/compute/models/cloudsync.go index fb3bf7ea37..11ae695fcc 100644 --- a/pkg/compute/models/cloudsync.go +++ b/pkg/compute/models/cloudsync.go @@ -1096,13 +1096,13 @@ func syncHostVMs(ctx context.Context, userCred mcclient.TokenCredential, syncRes return } - syncVMPeripherals(ctx, userCred, syncVMPairs[i].Local, syncVMPairs[i].Remote, localHost, provider, driver) + SyncVMPeripherals(ctx, userCred, syncVMPairs[i].Local, syncVMPairs[i].Remote, localHost, provider, driver) }() } } -func syncVMPeripherals( +func SyncVMPeripherals( ctx context.Context, userCred mcclient.TokenCredential, local *SGuest, diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 7f6ae78ca0..dc33771ee9 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -557,7 +557,7 @@ func (self *SGuest) StartMigrateTask( data.Set("guest_status", jsonutils.NewString(guestStatus)) dedicateMigrateTask := "GuestMigrateTask" - if self.GetHypervisor() != api.HYPERVISOR_KVM { + if len(self.ExternalId) > 0 { dedicateMigrateTask = "ManagedGuestMigrateTask" //托管私有云 } self.SetStatus(ctx, userCred, vmStatus, "") @@ -615,7 +615,7 @@ func (self *SGuest) StartGuestLiveMigrateTask( data.Set("guest_status", jsonutils.NewString(guestStatus)) dedicateMigrateTask := "GuestLiveMigrateTask" - if self.GetHypervisor() != api.HYPERVISOR_KVM { + if len(self.ExternalId) > 0 { dedicateMigrateTask = "ManagedGuestLiveMigrateTask" //托管私有云 } if task, err := taskman.TaskManager.NewTask(ctx, dedicateMigrateTask, self, userCred, data, parentTaskId, "", nil); err != nil { diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index be2a5a9d65..0d6868ee04 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -3068,7 +3068,7 @@ func (guest *SGuest) SyncAllWithCloudVM(ctx context.Context, userCred mcclient.T return errors.Wrap(err, "guest.syncWithCloudVM") } - syncVMPeripherals(ctx, userCred, guest, extVM, host, provider, driver) + SyncVMPeripherals(ctx, userCred, guest, extVM, host, provider, driver) return nil }