mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
fix: batch migrate in parrallel
This commit is contained in:
@@ -4444,7 +4444,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,
|
||||
@@ -4455,28 +4456,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
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -290,7 +290,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
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user