diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index f3c966b9ac..ce91c84351 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -4448,7 +4448,8 @@ func (manager *SGuestManager) PerformBatchMigrate(ctx context.Context, userCred guests[i] = *guest } - var hostGuests = map[string][]*api.GuestBatchMigrateParams{} + var hostGuests = map[string][]*SGuest{} + var hostGuestParams = map[string][]*api.GuestBatchMigrateParams{} for i := 0; i < len(guests); i++ { bmp := &api.GuestBatchMigrateParams{ Id: guests[i].Id, @@ -4459,28 +4460,34 @@ func (manager *SGuestManager) PerformBatchMigrate(ctx context.Context, userCred } guests[i].SetStatus(userCred, api.VM_START_MIGRATE, "batch migrate") if _, ok := hostGuests[guests[i].HostId]; ok { - hostGuests[guests[i].HostId] = append(hostGuests[guests[i].HostId], bmp) + hostGuests[guests[i].HostId] = append(hostGuests[guests[i].HostId], &guests[i]) + hostGuestParams[guests[i].HostId] = append(hostGuestParams[guests[i].HostId], bmp) } else { - hostGuests[guests[i].HostId] = []*api.GuestBatchMigrateParams{bmp} + hostGuests[guests[i].HostId] = []*SGuest{&guests[i]} + hostGuestParams[guests[i].HostId] = []*api.GuestBatchMigrateParams{bmp} } } - for hostId, params := range hostGuests { + for hostId, guests := range hostGuests { + params := hostGuestParams[hostId] kwargs := jsonutils.NewDict() kwargs.Set("guests", jsonutils.Marshal(params)) if len(preferHostId) > 0 { kwargs.Set("prefer_host_id", jsonutils.NewString(preferHostId)) } - host := HostManager.FetchHostById(hostId) - manager.StartHostGuestsMigrateTask(ctx, userCred, host, kwargs, "") + manager.StartHostGuestsMigrateTask(ctx, userCred, guests, kwargs, "") } return nil, nil } func (manager *SGuestManager) StartHostGuestsMigrateTask( ctx context.Context, userCred mcclient.TokenCredential, - host *SHost, kwargs *jsonutils.JSONDict, parentTaskId string, + guests []*SGuest, kwargs *jsonutils.JSONDict, parentTaskId string, ) error { - task, err := taskman.TaskManager.NewTask(ctx, "HostGuestsMigrateTask", host, userCred, kwargs, parentTaskId, "", nil) + taskItems := make([]db.IStandaloneModel, len(guests)) + for i := range guests { + taskItems[i] = guests[i] + } + task, err := taskman.TaskManager.NewParallelTask(ctx, "HostGuestsMigrateTask", taskItems, userCred, kwargs, parentTaskId, "", nil) if err != nil { log.Errorln(err) return err diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 269bf9f0d4..ffe336c078 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -5316,6 +5316,7 @@ func (host *SHost) MigrateSharedStorageServers(ctx context.Context, userCred mcc if err != nil { return errors.Wrapf(err, "host %s(%s) get guests", host.Name, host.Id) } + migGuests := []*SGuest{} hostGuests := []*api.GuestBatchMigrateParams{} for i := 0; i < len(guests); i++ { @@ -5333,11 +5334,12 @@ func (host *SHost) MigrateSharedStorageServers(ctx context.Context, userCred mcc } guests[i].SetStatus(userCred, api.VM_START_MIGRATE, "host down") hostGuests = append(hostGuests, bmp) + migGuests = append(migGuests, &guests[i]) } } kwargs := jsonutils.NewDict() kwargs.Set("guests", jsonutils.Marshal(hostGuests)) - return GuestManager.StartHostGuestsMigrateTask(ctx, userCred, host, kwargs, "") + return GuestManager.StartHostGuestsMigrateTask(ctx, userCred, migGuests, kwargs, "") } func (host *SHost) SetStatus(userCred mcclient.TokenCredential, status string, reason string) error { diff --git a/pkg/compute/tasks/guest_live_migrate_task.go b/pkg/compute/tasks/guest_live_migrate_task.go index 2b2d07ebce..5448827cfe 100644 --- a/pkg/compute/tasks/guest_live_migrate_task.go +++ b/pkg/compute/tasks/guest_live_migrate_task.go @@ -245,10 +245,26 @@ func (self *GuestMigrateTask) OnSrcPrepareComplete(ctx context.Context, guest *m 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) + err := jsonutils.NewDict() + err.Set("MigrateConfAndDiskFailedReason", data) + self.SetStage("OnUndeployTargetGuestSucc", err) + guest.StartUndeployGuestTask(ctx, self.UserCred, self.GetTaskId(), targetHostId) self.TaskFailed(ctx, guest, data) } +func (self *GuestMigrateTask) OnUndeployTargetGuestSucc(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + err, _ := self.Params.Get("MigrateConfAndDiskFailedReason") + self.TaskFailed(ctx, guest, err) +} + +func (self *GuestMigrateTask) OnUndeployTargetGuestSuccFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + prevErr, _ := self.Params.Get("MigrateConfAndDiskFailedReason") + err := jsonutils.NewDict() + err.Set("MigrateConfAndDiskFailedReason", prevErr) + err.Set("UndeployTargetGuestFailedReason", data) + self.TaskFailed(ctx, guest, err) +} + 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 == api.VM_RUNNING || guestStatus == api.VM_SUSPEND) { diff --git a/pkg/compute/tasks/host_guests_migrate_task.go b/pkg/compute/tasks/host_guests_migrate_task.go index 85bbcf7941..0d51164f39 100644 --- a/pkg/compute/tasks/host_guests_migrate_task.go +++ b/pkg/compute/tasks/host_guests_migrate_task.go @@ -34,7 +34,7 @@ type HostGuestsMigrateTask struct { taskman.STask } -func (self *HostGuestsMigrateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { +func (self *HostGuestsMigrateTask) OnInit(ctx context.Context, objs []db.IStandaloneModel, data jsonutils.JSONObject) { guests := make([]*api.GuestBatchMigrateParams, 0) err := self.Params.Unmarshal(&guests, "guests") if err != nil { @@ -43,51 +43,30 @@ func (self *HostGuestsMigrateTask) OnInit(ctx context.Context, obj db.IStandalon } preferHostId, _ := self.Params.GetString("prefer_host_id") - var guestMigrating bool - var migrateIndex int - for i := 0; i < len(guests); i++ { - guest := models.GuestManager.FetchGuestById(guests[i].Id) + self.SetStage("OnMigrateComplete", nil) + + for i := range objs { + guest := objs[i].(*models.SGuest) if guests[i].LiveMigrate { err := guest.StartGuestLiveMigrateTask( ctx, self.UserCred, guests[i].OldStatus, preferHostId, &guests[i].SkipCpuCheck, self.Id) if err != nil { log.Errorln(err) - continue - } else { - guestMigrating = true - migrateIndex = i - break } } else { err := guest.StartMigrateTask(ctx, self.UserCred, guests[i].RescueMode, false, guests[i].OldStatus, preferHostId, self.Id) if err != nil { log.Errorln(err) - continue - } else { - guestMigrating = true - migrateIndex = i - break } } } - if !guestMigrating { - if jsonutils.QueryBoolean(self.Params, "some_guest_migrate_failed", false) { - self.SetStageFailed(ctx, jsonutils.NewString("some guest migrate failed")) - } else { - self.SetStageComplete(ctx, nil) - } - } else { - guests := append(guests[:migrateIndex], guests[migrateIndex+1:]...) - params := jsonutils.NewDict() - params.Set("guests", jsonutils.Marshal(guests)) - self.SaveParams(params) - } } -func (self *HostGuestsMigrateTask) OnInitFailed(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { - kwargs := jsonutils.NewDict() - kwargs.Set("some_guest_migrate_failed", jsonutils.JSONTrue) - self.SaveParams(kwargs) - self.OnInit(ctx, obj, data) +func (self *HostGuestsMigrateTask) OnMigrateComplete(ctx context.Context, objs []db.IStandaloneModel, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} + +func (self *HostGuestsMigrateTask) OnMigrateCompleteFailed(ctx context.Context, objs []db.IStandaloneModel, data jsonutils.JSONObject) { + self.SetStageFailed(ctx, data) } diff --git a/pkg/compute/tasks/host_maintenance_task.go b/pkg/compute/tasks/host_maintenance_task.go index 3c5de73fbe..f4128dbf99 100644 --- a/pkg/compute/tasks/host_maintenance_task.go +++ b/pkg/compute/tasks/host_maintenance_task.go @@ -37,14 +37,36 @@ func init() { func (self *HostMaintainTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { host := obj.(*models.SHost) - guests, _ := self.Params.Get("guests") preferHostId, _ := self.Params.Get("prefer_host_id") + var hostGuests = []*api.GuestBatchMigrateParams{} + err := self.Params.Unmarshal(&hostGuests, "guests") + if err != nil { + self.TaskFailed(ctx, host, jsonutils.NewString(err.Error())) + return + } + + guests := make([]*models.SGuest, 0) + hostGuestParams := make([]*api.GuestBatchMigrateParams, 0) + for i := range hostGuests { + guest := models.GuestManager.FetchGuestById(hostGuests[i].Id) + if guest != nil { + guests = append(guests, guest) + hostGuestParams = append(hostGuestParams, hostGuests[i]) + } + } + + if len(guests) == 0 { + // no guest to migrate + self.SetStageComplete(ctx, nil) + return + } + kwargs := jsonutils.NewDict() - kwargs.Set("guests", guests) + kwargs.Set("guests", jsonutils.Marshal(hostGuestParams)) kwargs.Set("prefer_host_id", preferHostId) self.SetStage("OnGuestsMigrate", nil) - err := models.GuestManager.StartHostGuestsMigrateTask(ctx, self.UserCred, host, self.Params, self.Id) + err = models.GuestManager.StartHostGuestsMigrateTask(ctx, self.UserCred, guests, kwargs, self.Id) if err != nil { self.TaskFailed(ctx, host, jsonutils.NewString(err.Error())) return diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index b2e7eb8d39..2376ae5767 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -292,7 +292,11 @@ func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interfa if ctx != nil && len(appctx.AppContextTaskId(ctx)) >= 0 { hostutils.TaskFailed(ctx, fmt.Sprintf("Async start server failed: %s", err)) } - s.SyncStatus("") + needMigrate := jsonutils.QueryBoolean(data, "need_migrate", false) + // do not syncstatus if need_migrate + if !needMigrate { + s.SyncStatus("") + } return nil, err }