From 24c08ed178bd8b7411a0defc92dc7ab74ffc52fd Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Thu, 18 Oct 2018 16:22:27 +0800 Subject: [PATCH] server migrate all some bugfix --- cmd/climc/shell/servers.go | 27 ++ cmd/climc/shell/storages.go | 18 - pkg/cloudcommon/db/taskman/tasks.go | 16 +- pkg/compute/guestdrivers/kvm.go | 8 +- pkg/compute/models/disks.go | 4 +- pkg/compute/models/guestdisks.go | 3 +- pkg/compute/models/guests.go | 197 ++++++++- pkg/compute/tasks/guest_live_migrate_task.go | 403 +++++++++++++++++++ pkg/mcclient/options/servers.go | 13 + 9 files changed, 656 insertions(+), 33 deletions(-) create mode 100644 pkg/compute/tasks/guest_live_migrate_task.go diff --git a/cmd/climc/shell/servers.go b/cmd/climc/shell/servers.go index 0f947d663a..56a82adbac 100644 --- a/cmd/climc/shell/servers.go +++ b/cmd/climc/shell/servers.go @@ -4,6 +4,7 @@ import ( "fmt" "io/ioutil" + "yunion.io/x/jsonutils" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/modules" @@ -154,6 +155,32 @@ func init() { return nil }) + R(&options.ServerMigrateOptions{}, "server-migrate", "Migrate server", func(s *mcclient.ClientSession, opts *options.ServerMigrateOptions) error { + params, err := options.StructToParams(opts) + if err != nil { + return err + } + ret, err := modules.Servers.PerformAction(s, opts.ID, "migrate", params) + if err != nil { + return err + } + printObject(ret) + return nil + }) + + R(&options.ServerLiveMigrateOptions{}, "server-live-migrate", "Migrate server", func(s *mcclient.ClientSession, opts *options.ServerLiveMigrateOptions) error { + params, err := options.StructToParams(opts) + if err != nil { + return err + } + ret, err := modules.Servers.PerformAction(s, opts.ID, "live-migrate", params) + if err != nil { + return err + } + printObject(ret) + return nil + }) + R(&options.ServerResetOptions{}, "server-reset", "Reset servers", func(s *mcclient.ClientSession, opts *options.ServerResetOptions) error { params, err := options.StructToParams(opts) if err != nil { diff --git a/cmd/climc/shell/storages.go b/cmd/climc/shell/storages.go index 4d0cca8701..de96c3bf29 100644 --- a/cmd/climc/shell/storages.go +++ b/cmd/climc/shell/storages.go @@ -168,24 +168,6 @@ func init() { return nil }) - R(&StorageShowOptions{}, "storage-enable", "Enable a storage", func(s *mcclient.ClientSession, args *StorageShowOptions) error { - result, err := modules.Storages.PerformAction(s, args.ID, "enable", nil) - if err != nil { - return err - } - printObject(result) - return nil - }) - - R(&StorageShowOptions{}, "storage-disable", "Disable a storage", func(s *mcclient.ClientSession, args *StorageShowOptions) error { - result, err := modules.Storages.PerformAction(s, args.ID, "disable", nil) - if err != nil { - return err - } - printObject(result) - return nil - }) - type StorageCacheImageActionOptions struct { ID string `help:"ID or name of storage"` IMAGE string `help:"ID or name of image"` diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index ee9bdbd99a..9b4f175e46 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -286,11 +286,6 @@ func (manager *STaskManager) execTask(taskId string, data jsonutils.JSONObject) } log.Debugf("Do task %s(%s) with data %s at stage %s", taskType, taskId, data, baseTask.Stage) taskValue := reflect.New(taskType) - filled := reflectutils.FillEmbededStructValue(taskValue.Elem(), reflect.Indirect(reflect.ValueOf(baseTask))) - if !filled { - log.Errorf("Cannot locate baseTask embedded struct, give up...") - return - } if taskValue.Type().Implements(ITaskType) { execITask(taskValue, baseTask, data, false) } else if taskValue.Type().Implements(IBatchTaskType) { @@ -413,11 +408,18 @@ func execITask(taskValue reflect.Value, task *STask, odata jsonutils.JSONObject, params[2] = reflect.ValueOf(data) - log.Debugf("Call %s %s: %s with %s", task.TaskName, stageName, funcValue, params) + filled := reflectutils.FillEmbededStructValue(taskValue.Elem(), reflect.Indirect(reflect.ValueOf(task))) + if !filled { + log.Errorf("Cannot locate baseTask embedded struct, give up...") + return + } + log.Debugf("Call %s %s: %s with %s", task.TaskName, stageName, funcValue, params) funcValue.Call(params) - task.SaveRequestContext(&ctxData) + // call save request context + saveRequestContextFuncValue := taskValue.MethodByName("SaveRequestContext") + saveRequestContextFuncValue.Call([]reflect.Value{reflect.ValueOf(&ctxData)}) } func (task *STask) ScheduleRun(data jsonutils.JSONObject) { diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index d4108688e1..a045a8215a 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -180,7 +180,13 @@ func (self *SKVMGuestDriver) RequestUndeployGuestOnHost(ctx context.Context, gue header.Set("X-Auth-Token", task.GetUserCred().GetTokenString()) header.Set("X-Task-Id", task.GetTaskId()) header.Set("X-Region-Version", "v2") - _, res, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "DELETE", url, header, nil, false) + body := jsonutils.NewDict() + + // XXXXXXXX + if guest.HostId != host.Id { + body.Set("migrated", jsonutils.JSONTrue) + } + _, res, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "DELETE", url, header, body, false) if err != nil { return err } diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index b6f549bd1b..0af20288a7 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -1002,7 +1002,7 @@ type DiskInfo struct { MountPoint string Format string Size int64 - StorageId string + Storage string Backend string MediumType string Driver string @@ -1021,7 +1021,7 @@ func (self *SDisk) ToDiskInfo() DiskInfo { if storage == nil { return ret } - ret.StorageId = storage.Id + ret.Storage = storage.Id ret.Backend = storage.StorageType ret.MediumType = storage.MediumType return ret diff --git a/pkg/compute/models/guestdisks.go b/pkg/compute/models/guestdisks.go index 9060ad06f1..cb5834f328 100644 --- a/pkg/compute/models/guestdisks.go +++ b/pkg/compute/models/guestdisks.go @@ -150,8 +150,7 @@ func (self *SGuestdisk) GetJsonDescAtHost(host *SHost) jsonutils.JSONObject { desc.Add(jsonutils.NewString(disk.StorageId), "storage_id") localpath := disk.GetPathAtHost(host) if len(localpath) == 0 { - desc.Add(jsonutils.NewString(disk.GetFetchUrl()), "url") - storage := disk.GetStorage() + desc.Add(jsonutils.JSONTrue, "migrating") target := host.GetLeastUsedStorage(storage.StorageType) desc.Add(jsonutils.NewString(target.Id), "target_storage_id") disk.SetStatus(nil, DISK_START_MIGRATE, "migration") diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 66292bc030..1c7635018d 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -1829,6 +1829,129 @@ func (self *SGuest) CheckQemuVersion(qemuVer, compareVer string) bool { return true } +func (self *SGuest) AllowPerformMigrate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return self.IsOwner(userCred) +} + +func (self *SGuest) PerformMigrate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + if self.GetHypervisor() != HYPERVISOR_KVM { + return nil, httperrors.NewNotAcceptableError("Not allow for hypervisor %s", self.GetHypervisor()) + } + isRescueMode := jsonutils.QueryBoolean(data, "rescue_mode", false) + if !isRescueMode && self.Status != VM_READY { + return nil, httperrors.NewServerStatusError("Cannot normal migrate guest in status %s, try rescue mode or server-live-migrate?", self.Status) + } + if isRescueMode { + guestDisks := self.GetDisks() + for _, guestDisk := range guestDisks { + if utils.IsInStringArray( + guestDisk.GetDisk().GetStorage().StorageType, STORAGE_LOCAL_TYPES) { + return nil, httperrors.NewBadRequestError("Rescue mode requires all disk store in shared storages") + } + } + } + devices := self.GetIsolatedDevices() + if devices != nil && len(devices) > 0 { + return nil, httperrors.NewBadRequestError("Cannot migrate with isolated devices") + } + var preferHostId string + preferHost, _ := data.GetString("prefer_host") + if len(preferHost) > 0 { + if !userCred.IsSystemAdmin() { + return nil, httperrors.NewBadRequestError("Only system admin can assign host") + } + iHost, _ := HostManager.FetchByIdOrName(userCred, preferHost) + if iHost == nil { + return nil, httperrors.NewBadRequestError("Host %s not found", preferHost) + } + host := iHost.(*SHost) + preferHostId = host.Id + } + err := self.StartMigrateTask(ctx, userCred, isRescueMode, self.Status, preferHostId, "") + return nil, err +} + +func (self *SGuest) StartMigrateTask(ctx context.Context, userCred mcclient.TokenCredential, isRescueMode bool, guestStatus, preferHostId, parentTaskId string) error { + data := jsonutils.NewDict() + if isRescueMode { + data.Set("is_rescue_mode", jsonutils.JSONTrue) + } + if len(preferHostId) > 0 { + data.Set("prefer_host_id", jsonutils.NewString(preferHostId)) + } + data.Set("guest_status", jsonutils.NewString(guestStatus)) + if task, err := taskman.TaskManager.NewTask(ctx, "GuestMigrateTask", self, userCred, data, parentTaskId, "", nil); err != nil { + log.Errorf(err.Error()) + return err + } else { + task.ScheduleRun(nil) + } + return nil +} + +func (self *SGuest) AllowPerformLiveMigrate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return self.IsOwner(userCred) +} + +func (self *SGuest) PerformLiveMigrate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + if self.GetHypervisor() != HYPERVISOR_KVM { + return nil, httperrors.NewNotAcceptableError("Not allow for hypervisor %s", self.GetHypervisor()) + } + imageId := self.GetDisks()[0].GetDisk().TemplateId + image, err := CachedimageManager.GetImageById(ctx, userCred, imageId, false) + if err != nil { + return nil, err + } + if image.DiskFormat != "qcow2" { + return nil, httperrors.NewBadRequestError("Live migrate only support image fromat qocw2") + } + if utils.IsInStringArray(self.Status, []string{VM_RUNNING, VM_SUSPEND}) { + cdrom := self.getCdrom() + if cdrom != nil && len(cdrom.ImageId) > 0 { + return nil, httperrors.NewBadRequestError("Cannot migrate with cdrom") + } + devices := self.GetIsolatedDevices() + if devices != nil && len(devices) > 0 { + return nil, httperrors.NewBadRequestError("Cannot migrate with isolated devices") + } + if !self.CheckQemuVersion(self.GetQemuVersion(userCred), "1.1.2") { + return nil, httperrors.NewBadRequestError("Cannot do live migrate, too low qemu version") + } + var preferHostId string + preferHost, _ := data.GetString("prefer_host") + if len(preferHost) > 0 { + if !userCred.IsSystemAdmin() { + return nil, httperrors.NewBadRequestError("Only system admin can assign host") + } + iHost, _ := HostManager.FetchByIdOrName(userCred, preferHost) + if iHost == nil { + return nil, httperrors.NewBadRequestError("Host %s not found", preferHost) + } + host := iHost.(*SHost) + preferHostId = host.Id + } + err := self.StartGuestLiveMigrateTask(ctx, userCred, self.Status, preferHostId, "") + return nil, err + } + return nil, httperrors.NewBadRequestError("Cannot live migrate in status %s", self.Status) +} + +func (self *SGuest) StartGuestLiveMigrateTask(ctx context.Context, userCred mcclient.TokenCredential, guestStatus, preferHostId, parentTaskId string) error { + self.SetStatus(userCred, VM_START_MIGRATE, "") + data := jsonutils.NewDict() + if len(preferHostId) > 0 { + data.Set("prefer_host_id", jsonutils.NewString(preferHostId)) + } + data.Set("guest_status", jsonutils.NewString(guestStatus)) + if task, err := taskman.TaskManager.NewTask(ctx, "GuestLiveMigrateTask", self, userCred, data, parentTaskId, "", nil); err != nil { + log.Errorf(err.Error()) + return err + } else { + task.ScheduleRun(nil) + } + return nil +} + func (self *SGuest) AllowPerformDeploy(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { return self.IsOwner(userCred) } @@ -3974,18 +4097,18 @@ func (self *SGuest) AllowPerformSuspend(ctx context.Context, userCred mcclient.T func (self *SGuest) PerformSuspend(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { if self.Status == VM_RUNNING { - err := self.StartSuspendTask(ctx, userCred) + err := self.StartSuspendTask(ctx, userCred, "") return nil, err } return nil, httperrors.NewInvalidStatusError("Cannot suspend VM in status %s", self.Status) } -func (self *SGuest) StartSuspendTask(ctx context.Context, userCred mcclient.TokenCredential) error { +func (self *SGuest) StartSuspendTask(ctx context.Context, userCred mcclient.TokenCredential, parentTaskId string) error { err := self.SetStatus(userCred, VM_SUSPEND, "do suspend") if err != nil { return err } - return self.GetDriver().StartSuspendTask(ctx, userCred, self, nil, "") + return self.GetDriver().StartSuspendTask(ctx, userCred, self, nil, parentTaskId) } func (self *SGuest) AllowPerformStart(ctx context.Context, @@ -4702,3 +4825,71 @@ func (self *SGuest) getSchedDesc() jsonutils.JSONObject { return desc } + +func (self *SGuest) GetApptags() []string { + tagsStr := self.GetMetadata("app_tags", nil) + if len(tagsStr) > 0 { + return strings.Split(tagsStr, ",") + } + return nil +} + +func (self *SGuest) ToSchedDesc() *jsonutils.JSONDict { + desc := jsonutils.NewDict() + desc.Set("id", jsonutils.NewString(self.Id)) + desc.Set("name", jsonutils.NewString(self.Name)) + desc.Set("vmem_size", jsonutils.NewInt(int64(self.VmemSize))) + desc.Set("vcpu_count", jsonutils.NewInt(int64(self.VcpuCount))) + self.FillGroupSchedDesc(desc) + self.FillDiskSchedDesc(desc) + self.FillNetSchedDesc(desc) + if len(self.HostId) > 0 && regutils.MatchUUID(self.HostId) { + desc.Set("host_id", jsonutils.NewString(self.HostId)) + } + desc.Set("owner_tenant_id", jsonutils.NewString(self.ProjectId)) + tags := self.GetApptags() + for i := 0; i < len(tags); i++ { + desc.Set(tags[i], jsonutils.JSONTrue) + } + desc.Set("hypervisor", jsonutils.NewString(self.GetHypervisor())) + return desc +} + +func (self *SGuest) FillGroupSchedDesc(desc *jsonutils.JSONDict) { + groups := make([]SGroupguest, 0) + err := GroupguestManager.Query().Equals("guest_id", self.Id).All(&groups) + if err != nil { + log.Errorln(err) + return + } + for i := 0; i < len(groups); i++ { + desc.Set(fmt.Sprintf("srvtag.%d", i), + jsonutils.NewString(fmt.Sprintf("%s:%s", groups[i].SrvtagId, groups[i].Tag))) + } +} + +func (self *SGuest) FillDiskSchedDesc(desc *jsonutils.JSONDict) { + guestDisks := make([]SGuestdisk, 0) + err := GuestdiskManager.Query().Equals("guest_id", self.Id).All(&guestDisks) + if err != nil { + log.Errorln(err) + return + } + for i := 0; i < len(guestDisks); i++ { + desc.Set(fmt.Sprintf("disk.%d", i), jsonutils.Marshal(guestDisks[i].ToDiskInfo())) + } +} + +func (self *SGuest) FillNetSchedDesc(desc *jsonutils.JSONDict) { + guestNetworks := make([]SGuestnetwork, 0) + err := GuestnetworkManager.Query().Equals("guest_id", self.Id).All(&guestNetworks) + if err != nil { + log.Errorln(err) + return + } + for i := 0; i < len(guestNetworks); i++ { + desc.Set(fmt.Sprintf("net.%d", i), + jsonutils.NewString(fmt.Sprintf("%s:%s", + guestNetworks[i].NetworkId, guestNetworks[i].IpAddr))) + } +} diff --git a/pkg/compute/tasks/guest_live_migrate_task.go b/pkg/compute/tasks/guest_live_migrate_task.go new file mode 100644 index 0000000000..7a36a4255b --- /dev/null +++ b/pkg/compute/tasks/guest_live_migrate_task.go @@ -0,0 +1,403 @@ +package tasks + +import ( + "context" + "fmt" + "net/http" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/utils" + + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "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/compute/options" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/util/httputils" +) + +type GuestMigrateTask struct { + SGuestBaseTask +} + +type GuestLiveMigrateTask struct { + GuestMigrateTask +} + +func init() { + taskman.RegisterTask(GuestLiveMigrateTask{}) + taskman.RegisterTask(GuestMigrateTask{}) +} + +func (self *GuestMigrateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + guest.SetStatus(self.UserCred, models.VM_MIGRATING, "") + db.OpsLog.LogEvent(guest, db.ACT_MIGRATING, "", self.UserCred) + self.FindCandidataTarget(ctx, guest) +} + +func (self *GuestMigrateTask) FindCandidataTarget(ctx context.Context, guest *models.SGuest) { + schedDesc := guest.ToSchedDesc() + if self.Params.Contains("prefer_host_id") { + preferHostId, _ := self.Params.Get("prefer_host_id") + schedDesc.Set("prefer_host_id", preferHostId) + } + s := auth.GetAdminSession(options.Options.Region, "") + results, err := modules.SchedManager.DoSchedule(s, schedDesc, 1) + if err != nil { + self.TaskFailed(ctx, guest, fmt.Sprintf("Do schedule error %s", err)) + } else { + self.OnScheduleComplete(ctx, guest, results) + } +} + +func (self *GuestMigrateTask) OnScheduleComplete(ctx context.Context, guest *models.SGuest, results []jsonutils.JSONObject) { + if len(results) != 1 { + self.TaskFailed(ctx, guest, "Schedule failed") + return + } + var targetHostId string + if results[0].Contains("candidate") { + targetHostId, _ = results[0].GetString("candidate", "id") + } else if results[0].Contains("error") { + msg, _ := results[0].Get("error") + self.TaskFailed(ctx, guest, msg.String()) + return + } else { + msg := fmt.Sprintf("Unknown scheduler result %s", results[0]) + self.TaskFailed(ctx, guest, msg) + return + } + + targetHost := models.HostManager.FetchHostById(targetHostId) + if targetHost == nil { + self.TaskFailed(ctx, guest, "target host not found?") + return + } + db.OpsLog.LogEvent(guest, db.ACT_MIGRATING, fmt.Sprintf("guest start migrate from host %s to %s", guest.HostId, targetHostId), self.UserCred) + + body := jsonutils.NewDict() + body.Set("target_host_id", jsonutils.NewString(targetHostId)) + + disks := guest.GetDisks() + disk := disks[0].GetDisk() + isLocalStorage := utils.IsInStringArray(disk.GetStorage().StorageType, + models.STORAGE_LOCAL_TYPES) + if isLocalStorage { + body.Set("is_local_storage", jsonutils.JSONTrue) + } else { + body.Set("is_local_storage", jsonutils.JSONFalse) + } + + self.SetStage("OnCachedImageComplete", body) + // prepare disk for migration + if isLocalStorage { + targetStorageCache := targetHost.GetLocalStoragecache() + if targetStorageCache != nil { + targetStorageCache.StartImageCacheTask(ctx, self.UserCred, disk.TemplateId, false, self.GetTaskId()) + } + } else { + self.OnSrcPrepareComplete(ctx, guest, nil) + } +} + +// For local storage get disk info +func (self *GuestMigrateTask) OnCachedImageComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + header := http.Header{} + header.Set("X-Auth-Token", self.GetUserCred().GetTokenString()) + header.Set("X-Task-Id", self.GetTaskId()) + header.Set("X-Region-Version", "v2") + body := jsonutils.NewDict() + guestStatus, _ := self.Params.GetString("guest_status") + if !jsonutils.QueryBoolean(self.Params, "is_rescue_mode", false) && (guestStatus == models.VM_RUNNING || guestStatus == models.VM_SUSPEND) { + body.Set("live_migrate", jsonutils.JSONTrue) + } + + host := guest.GetHost() + url := fmt.Sprintf("%s/servers/%s/src-prepare-migrate", host.ManagerUri, guest.Id) + self.SetStage("OnSrcPrepareComplete", nil) + _, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", + url, header, body, false) + if err != nil { + self.TaskFailed(ctx, guest, fmt.Sprintf("Prepare migrage failed: %s", err)) + return + } +} + +func (self *GuestMigrateTask) OnSrcPrepareCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.TaskFailed(ctx, guest, data.String()) +} + +func (self *GuestMigrateTask) OnSrcPrepareComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + targetHostId, _ := self.Params.GetString("target_host_id") + targetHost := models.HostManager.FetchHostById(targetHostId) + var body *jsonutils.JSONDict + var hasError bool + if jsonutils.QueryBoolean(self.Params, "is_local_storage", false) { + body, hasError = self.localStorageMigrateConf(ctx, guest, targetHost, data) + } else { + body, hasError = self.sharedStorageMigrateConf(ctx, guest, targetHost) + } + if hasError { + return + } + guestStatus, _ := self.Params.GetString("guest_status") + if !jsonutils.QueryBoolean(self.Params, "is_rescue_mode", false) && (guestStatus == models.VM_RUNNING || guestStatus == models.VM_SUSPEND) { + body.Set("live_migrate", jsonutils.JSONTrue) + } + + headers := http.Header{} + headers.Set("X-Auth-Token", self.GetUserCred().GetTokenString()) + headers.Set("X-Task-Id", self.GetTaskId()) + headers.Set("X-Region-Version", "v2") + + url := fmt.Sprintf("%s/servers/%s/dest-prepare-migrate", targetHost.ManagerUri, guest.Id) + self.SetStage("OnMigrateConfAndDiskComplete", nil) + _, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), + ctx, "POST", url, headers, body, false) + if err != nil { + self.TaskFailed(ctx, guest, err.Error()) + } +} + +func (self *GuestMigrateTask) OnMigrateConfAndDiskCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + targetHostId, _ := self.Params.GetString("target_host_id") + guest.StartUndeployGuestTask(ctx, self.UserCred, "", targetHostId) + self.TaskFailed(ctx, guest, data.String()) +} + +func (self *GuestMigrateTask) OnMigrateConfAndDiskComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guestStatus, _ := self.Params.GetString("guest_status") + if !jsonutils.QueryBoolean(self.Params, "is_rescue_mode", false) && (guestStatus == models.VM_RUNNING || guestStatus == models.VM_SUSPEND) { + // Live migrate + self.SetStage("OnStartDestComplete", nil) + } else { + // Normal migrate + self.OnNormalMigrateComplete(ctx, guest, data) + } +} + +func (self *GuestMigrateTask) OnNormalMigrateComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + oldHostId := guest.HostId + self.setGuest(ctx, guest) + guestStatus, _ := self.Params.GetString("guest_status") + guest.SetStatus(self.UserCred, guestStatus, "") + if jsonutils.QueryBoolean(self.Params, "is_rescue_mode", false) { + guest.StartGueststartTask(ctx, self.UserCred, nil, "") + } + self.SetStage("OnUndeployOldHostSucc", nil) + guest.StartUndeployGuestTask(ctx, self.UserCred, self.GetTaskId(), oldHostId) +} + +func (self *GuestMigrateTask) OnUndeployOldHostSucc(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} + +func (self *GuestMigrateTask) sharedStorageMigrateConf(ctx context.Context, guest *models.SGuest, targetHost *models.SHost) (*jsonutils.JSONDict, bool) { + body := jsonutils.NewDict() + body.Set("is_local_storage", jsonutils.JSONFalse) + body.Set("qemu_version", jsonutils.NewString(guest.GetQemuVersion(self.UserCred))) + targetDesc := guest.GetJsonDescAtHypervisor(ctx, targetHost) + body.Set("desc", targetDesc) + return body, false +} + +func (self *GuestMigrateTask) localStorageMigrateConf(ctx context.Context, + guest *models.SGuest, targetHost *models.SHost, data jsonutils.JSONObject) (*jsonutils.JSONDict, bool) { + body := jsonutils.NewDict() + if data != nil { + body.Update(data.(*jsonutils.JSONDict)) + } + params := jsonutils.NewDict() + disks := guest.GetDisks() + for i := 0; i < len(disks); i++ { + snapshots := models.SnapshotManager.GetDiskSnapshots(disks[i].DiskId) + snapshotIds := jsonutils.NewArray() + for j := 0; j < len(snapshots); j++ { + snapshotIds.Add(jsonutils.NewString(snapshots[j].Id)) + } + params.Set(disks[i].DiskId, snapshotIds) + } + + sourceHost := guest.GetHost() + snapshotsUri := fmt.Sprintf("%s/download/snapshots/", sourceHost.ManagerUri) + disksUri := fmt.Sprintf("%s/download/disks/", sourceHost.ManagerUri) + serverUrl := fmt.Sprintf("%s/download/servers/%s", sourceHost.ManagerUri, guest.Id) + + body.Set("src_snapshots", params) + body.Set("snapshots_uri", jsonutils.NewString(snapshotsUri)) + body.Set("disks_uri", jsonutils.NewString(disksUri)) + body.Set("server_url", jsonutils.NewString(serverUrl)) + body.Set("qemu_version", jsonutils.NewString(guest.GetQemuVersion(self.UserCred))) + targetDesc := guest.GetJsonDescAtHypervisor(ctx, targetHost) + jsonDisks, _ := targetDesc.Get("disks") + if jsonDisks == nil { + self.TaskFailed(ctx, guest, "Get jsonDisks error") + return nil, true + } + disksDesc, _ := jsonDisks.GetArray() + if len(disksDesc) == 0 { + self.TaskFailed(ctx, guest, "Get disksDesc error") + return nil, true + } + targetStorageId, _ := disksDesc[0].GetString("target_storage_id") + if len(targetStorageId) == 0 { + self.TaskFailed(ctx, guest, "Get targetStorageId error") + return nil, true + } + + targetStorage := targetHost.GetHoststorageOfId(targetStorageId) + sourceStorage := sourceHost.GetHoststorageOfId(disks[0].GetDisk().StorageId) + if sourceStorage.MountPoint != targetStorage.MountPoint { + self.TaskFailed(ctx, guest, fmt.Sprintf("target host %s storage"+ + "mount point is different with source storage", targetHost.Id)) + return nil, true + } + body.Set("desc", targetDesc) + body.Set("is_local_storage", jsonutils.JSONTrue) + return body, false +} + +func (self *GuestLiveMigrateTask) OnStartDestComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + liveMigrateDestPort, err := data.Get("live_migrate_dest_port") + if err != nil { + self.TaskFailed(ctx, guest, fmt.Sprintf("Get migrate port error: %s", err)) + return + } + + targetHostId, _ := self.Params.GetString("target_host_id") + targetHost := models.HostManager.FetchHostById(targetHostId) + + body := jsonutils.NewDict() + isLocalStorage, _ := self.Params.Get("is_local_storage") + body.Set("is_local_storage", isLocalStorage) + body.Set("live_migrate_dest_port", liveMigrateDestPort) + body.Set("dest_ip", jsonutils.NewString(targetHost.AccessIp)) + + headers := http.Header{} + headers.Set("X-Auth-Token", self.GetUserCred().GetTokenString()) + headers.Set("X-Task-Id", self.GetTaskId()) + headers.Set("X-Region-Version", "v2") + + host := guest.GetHost() + url := fmt.Sprintf("%s/servers/%s/live-migrate", host.ManagerUri, guest.Id) + self.SetStage("OnLiveMigrateComplete", nil) + _, _, err = httputils.JSONRequest(httputils.GetDefaultClient(), + ctx, "POST", url, headers, body, false) + if err != nil { + self.OnLiveMigrateCompleteFailed(ctx, guest, jsonutils.NewString(err.Error())) + } +} + +func (self *GuestLiveMigrateTask) OnStartDestCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + targetHostId, _ := self.Params.GetString("target_host_id") + guest.StartUndeployGuestTask(ctx, self.UserCred, "", targetHostId) + self.TaskFailed(ctx, guest, data.String()) +} + +func (self *GuestMigrateTask) setGuest(ctx context.Context, guest *models.SGuest) error { + targetHostId, _ := self.Params.GetString("target_host_id") + if jsonutils.QueryBoolean(self.Params, "is_local_storage", false) { + targetHost := models.HostManager.FetchHostById(targetHostId) + targetStorage := targetHost.GetLeastUsedStorage(models.STORAGE_LOCAL) + guestDisks := guest.GetDisks() + for i := 0; i < len(guestDisks); i++ { + disk := guestDisks[i].GetDisk() + disk.GetModelManager().TableSpec().Update(disk, func() error { + disk.Status = models.DISK_READY + disk.StorageId = targetStorage.Id + return nil + }) + snapshots := models.SnapshotManager.GetDiskSnapshots(disk.Id) + for _, snapshot := range snapshots { + snapshot.GetModelManager().TableSpec().Update(snapshot, func() error { + snapshot.StorageId = targetStorage.Id + return nil + }) + } + } + } + oldHost := guest.GetHost() + oldHost.ClearSchedDescCache() + err := guest.SetHostId(targetHostId) + if err != nil { + return err + } + return nil +} + +func (self *GuestLiveMigrateTask) OnLiveMigrateCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + targetHostId, _ := self.Params.GetString("target_host_id") + guest.StartUndeployGuestTask(ctx, self.UserCred, "", targetHostId) + self.TaskFailed(ctx, guest, data.String()) +} + +func (self *GuestLiveMigrateTask) OnLiveMigrateComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + headers := http.Header{} + headers.Set("X-Auth-Token", self.GetUserCred().GetTokenString()) + headers.Set("X-Task-Id", self.GetTaskId()) + headers.Set("X-Region-Version", "v2") + body := jsonutils.NewDict() + body.Set("live_migrate", jsonutils.JSONTrue) + targetHostId, _ := self.Params.GetString("target_host_id") + + self.SetStage("OnResumeDestGuestComplete", nil) + targetHost := models.HostManager.FetchHostById(targetHostId) + url := fmt.Sprintf("%s/servers/%s/resume", targetHost.ManagerUri, guest.Id) + _, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), + ctx, "POST", url, headers, body, false) + if err != nil { + self.TaskFailed(ctx, guest, err.Error()) + } +} + +func (self *GuestLiveMigrateTask) OnResumeDestGuestCompleteFailed(ctx context.Context, + guest *models.SGuest, data jsonutils.JSONObject) { + targetHostId, _ := self.Params.GetString("target_host_id") + + guest.StartUndeployGuestTask(ctx, self.UserCred, "", targetHostId) + self.TaskFailed(ctx, guest, data.String()) +} + +func (self *GuestLiveMigrateTask) OnResumeDestGuestComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + oldHostId := guest.HostId + err := self.setGuest(ctx, guest) + if err != nil { + self.TaskFailed(ctx, guest, err.Error()) + } + self.SetStage("OnUndeploySrcGuestComplete", nil) + err = guest.StartUndeployGuestTask(ctx, self.UserCred, self.GetTaskId(), oldHostId) + if err != nil { + self.TaskFailed(ctx, guest, err.Error()) + } +} + +func (self *GuestLiveMigrateTask) OnUndeploySrcGuestComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + db.OpsLog.LogEvent(guest, db.ACT_MIGRATE, "", self.UserCred) + status, _ := self.Params.GetString("guest_status") + if status != models.VM_RUNNING { + guest.SetStatus(self.UserCred, status, "") + self.SetStageComplete(ctx, nil) + } else { + self.SetStage("OnStartGeustComplete", nil) + guest.StartGueststartTask(ctx, self.UserCred, nil, self.GetTaskId()) + } +} + +func (self *GuestLiveMigrateTask) OnStartGeustComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} + +func (self *GuestMigrateTask) TaskFailed(ctx context.Context, guest *models.SGuest, reason string) { + status, _ := self.Params.GetString("guest_status") + if status != models.VM_RUNNING { + guest.SetStatus(self.UserCred, status, "") + } else { + guest.StartGueststartTask(ctx, self.UserCred, nil, "") + } + db.OpsLog.LogEvent(guest, db.ACT_MIGRATE_FAIL, reason, self.UserCred) + self.SetStageFailed(ctx, reason) + notifyclient.NotifySystemError(guest.Id, guest.Name, models.VM_MIGRATE_FAILED, reason) +} diff --git a/pkg/mcclient/options/servers.go b/pkg/mcclient/options/servers.go index 932ebf6f04..30c9f17c76 100644 --- a/pkg/mcclient/options/servers.go +++ b/pkg/mcclient/options/servers.go @@ -294,3 +294,16 @@ type ServerRestartOptions struct { ID []string `help:"ID of servers to operate" metavar:"SERVER" json:"-"` IsForce *bool `help:"Force reset or not; default false" json:"is_force"` } + +type ServerMigrateOptions struct { + ID string `help:"ID of server" json:"-"` + PreferHost string `help:"Server migration prefer host id or name" json:"prefer_host"` + RescueMode *bool `help:"Migrate server in rescue mode, + all disk must store in shared storage; + default false" json:"rescue_mode"` +} + +type ServerLiveMigrateOptions struct { + ID string `help:"ID of server" json:"-"` + PreferHost string `help:"Server migration prefer host id or name" json:"prefer_host"` +}