Merge pull request #9441 from wanyaoqi/bugfix/wyq/fix-change-config

bugfix(region): server change config rebase to sched task
This commit is contained in:
Zexi Li
2020-12-17 09:50:42 +08:00
committed by GitHub
3 changed files with 128 additions and 69 deletions
+11 -52
View File
@@ -2449,7 +2449,6 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
disks := self.GetDisks()
var addDisk int
var newDiskIdx = 0
var diskSizes = make(map[string]int, 0)
var newDisks = make([]*api.DiskConfig, 0)
var resizeDisks = jsonutils.NewArray()
@@ -2460,6 +2459,7 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
}
}
var schedInputDisks = make([]*api.DiskConfig, 0)
var diskIdx = 1
for _, diskConf := range inputDisks {
diskConf, err = parseDiskInfo(ctx, userCred, diskConf)
@@ -2471,20 +2471,10 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
}
if diskConf.SizeMb > 0 {
if diskIdx >= len(disks) {
// 这里backeend为空时,qcloud有可能会选择local_ssd作为后端存储,会导致报错(主要是climc)
storage := host.GetLeastUsedStorage(diskConf.Backend)
if storage == nil {
return nil, httperrors.NewResourceNotReadyError("host not connect storage %s", diskConf.Backend)
}
_, ok := diskSizes[storage.Id]
if !ok {
diskSizes[storage.Id] = 0
}
diskSizes[storage.Id] = diskSizes[storage.Id] + diskConf.SizeMb
diskConf.Storage = storage.Id
newDisks = append(newDisks, diskConf)
newDiskIdx += 1
addDisk += diskConf.SizeMb
schedInputDisks = append(schedInputDisks, diskConf)
} else {
disk := disks[diskIdx].GetDisk()
oldSize := disk.DiskSize
@@ -2495,40 +2485,17 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
resizeDisks.Add(arr)
addDisk += diskConf.SizeMb - oldSize
storage := disks[diskIdx].GetDisk().GetStorage()
_, ok := diskSizes[storage.Id]
if !ok {
diskSizes[storage.Id] = 0
}
err = self.ValidateResizeDisk(disk, storage)
if err != nil {
return nil, httperrors.NewUnsupportOperationError("%v", err)
}
if !storage.IsEmulated && storage.GetFreeCapacity() < int64(addDisk) {
return nil, httperrors.NewInsufficientResourceError("Not enough free space")
}
diskSizes[storage.Id] = diskSizes[storage.Id] + diskConf.SizeMb - oldSize
schedInputDisks = append(schedInputDisks, &api.DiskConfig{
SizeMb: addDisk,
Index: diskConf.Index,
Storage: storage.Id,
})
}
}
}
diskIdx += 1
}
provider, e := self.GetHost().GetProviderFactory()
if e != nil || !provider.IsPublicCloud() {
for storageId, needSize := range diskSizes {
iStorage, err := StorageManager.FetchById(storageId)
if err != nil {
return nil, httperrors.NewBadRequestError("Fetch storage error: %s", err)
}
storage := iStorage.(*SStorage)
if !storage.IsEmulated && storage.GetFreeCapacity() < int64(needSize) {
return nil, httperrors.NewInsufficientResourceError("Not enough free space")
}
}
} else {
log.Debugf("Skip storage free capacity validating for public cloud: %s", provider.GetId())
}
if resizeDisks.Length() > 0 {
confs.Add(resizeDisks, "resize")
}
@@ -2545,7 +2512,8 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
}
// schedulr forecast
schedDesc := self.changeConfToSchedDesc(addCpu, addMem, addDisk)
schedDesc := self.changeConfToSchedDesc(addCpu, addMem, schedInputDisks)
confs.Set("sched_desc", jsonutils.Marshal(schedDesc))
s := auth.GetAdminSession(ctx, options.Options.Region, "")
canChangeConf, err := modules.SchedManager.DoScheduleForecast(s, schedDesc, 1)
if err != nil {
@@ -2580,28 +2548,19 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
}
if len(newDisks) > 0 {
err := self.CreateDisksOnHost(ctx, userCred, host, newDisks, pendingUsage, false, false, nil, nil, false)
if err != nil {
quotas.CancelPendingUsage(ctx, userCred, pendingUsage, pendingUsage, false)
return nil, httperrors.NewBadRequestError("Create disk on host error: %s", err)
}
confs.Add(jsonutils.Marshal(newDisks), "create")
}
self.StartChangeConfigTask(ctx, userCred, confs, "", pendingUsage)
return nil, nil
}
func (self *SGuest) changeConfToSchedDesc(addCpu, addMem, addDisk int) *schedapi.ScheduleInput {
guestDisks := self.GetDisks()
diskInfo := guestDisks[0].ToDiskConfig()
diskInfo.SizeMb = addDisk
func (self *SGuest) changeConfToSchedDesc(addCpu, addMem int, schedInputDisks []*api.DiskConfig) *schedapi.ScheduleInput {
desc := &schedapi.ScheduleInput{
ServerConfig: schedapi.ServerConfig{
ServerConfigs: &api.ServerConfigs{
Hypervisor: self.Hypervisor,
PreferHost: self.HostId,
Disks: []*api.DiskConfig{diskInfo},
Disks: schedInputDisks,
},
Memory: addMem,
Ncpu: addCpu,
+88 -3
View File
@@ -22,6 +22,7 @@ import (
"yunion.io/x/log"
api "yunion.io/x/onecloud/pkg/apis/compute"
schedapi "yunion.io/x/onecloud/pkg/apis/scheduler"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
"yunion.io/x/onecloud/pkg/cloudcommon/db/quotas"
@@ -32,7 +33,7 @@ import (
)
type GuestChangeConfigTask struct {
SGuestBaseTask
SSchedTask
}
func init() {
@@ -40,12 +41,69 @@ func init() {
}
func (self *GuestChangeConfigTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
StartScheduleObjects(ctx, self, nil)
}
func (self *GuestChangeConfigTask) GetSchedParams() (*schedapi.ScheduleInput, error) {
schedInput := new(schedapi.ScheduleInput)
err := self.Params.Unmarshal(schedInput, "sched_desc")
if err != nil {
return nil, err
}
return schedInput, nil
}
func (self *GuestChangeConfigTask) OnStartSchedule(obj IScheduleModel) {
// do nothing
}
func (self *GuestChangeConfigTask) OnScheduleFailCallback(ctx context.Context, obj IScheduleModel, reason jsonutils.JSONObject) {
// do nothing
}
func (self *GuestChangeConfigTask) OnScheduleFailed(ctx context.Context, reason jsonutils.JSONObject) {
obj := self.GetObject()
guest := obj.(*models.SGuest)
self.markStageFailed(ctx, guest, reason)
}
func (self *GuestChangeConfigTask) SaveScheduleResult(ctx context.Context, obj IScheduleModel, target *schedapi.CandidateResource) {
// must get object from task, because of obj is nil
guest := self.GetObject().(*models.SGuest)
self.Params.Set("sched_session_id", jsonutils.NewString(target.SessionId))
if self.Params.Contains("create") {
disks := make([]*api.DiskConfig, 0)
err := self.Params.Unmarshal(&disks, "create")
if err != nil {
self.markStageFailed(ctx, guest, jsonutils.NewString(err.Error()))
return
}
var resizeDisksCount = 0
if self.Params.Contains("resize") {
iResizeDisks, err := self.Params.Get("resize")
if err != nil {
self.markStageFailed(ctx, guest, jsonutils.NewString(err.Error()))
return
}
resizeDisksCount = iResizeDisks.(*jsonutils.JSONArray).Length()
}
for i := 0; i < len(disks); i++ {
disks[i].Storage = target.Disks[resizeDisksCount+i].StorageIds[0]
}
self.Params.Set("create", jsonutils.Marshal(disks))
}
self.SetStage("StartResizeDisks", nil)
self.StartResizeDisks(ctx, guest, nil)
}
func (self *GuestChangeConfigTask) StartResizeDisks(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) {
_, err := self.Params.Get("resize")
if err == nil {
self.SetStage("OnDisksResizeComplete", nil)
self.OnDisksResizeComplete(ctx, obj, data)
self.OnDisksResizeComplete(ctx, guest, data)
} else {
guest := obj.(*models.SGuest)
self.DoCreateDisksTask(ctx, guest)
}
}
@@ -118,6 +176,11 @@ func (self *GuestChangeConfigTask) DoCreateDisksTask(ctx context.Context, guest
self.OnCreateDisksComplete(ctx, guest, nil)
return
}
err = guest.CreateDisksOnHost(ctx, self.UserCred, guest.GetHost(), disks, nil, false, false, nil, nil, false)
if err != nil {
self.markStageFailed(ctx, guest, jsonutils.NewString(err.Error()))
return
}
self.SetStage("OnCreateDisksComplete", nil)
guest.StartGuestCreateDiskTask(ctx, self.UserCred, disks, self.GetTaskId())
}
@@ -322,3 +385,25 @@ func (self *GuestChangeConfigTask) markStageFailed(ctx context.Context, guest *m
notifyclient.NotifyError(ctx, self.UserCred, guest.GetId(), guest.GetName(), logclient.ACT_VM_CHANGE_FLAVOR, reason.String())
self.SetStageFailed(ctx, reason)
}
func (self *GuestChangeConfigTask) SetStageFailed(ctx context.Context, reason jsonutils.JSONObject) {
guest := self.GetObject().(*models.SGuest)
hostId := guest.HostId
sessionId, _ := self.Params.GetString("sched_session_id")
lockman.LockRawObject(ctx, models.HostManager.KeywordPlural(), hostId)
defer lockman.ReleaseRawObject(ctx, models.HostManager.KeywordPlural(), hostId)
models.HostManager.ClearSchedDescSessionCache(hostId, sessionId)
self.SSchedTask.SetStageFailed(ctx, reason)
}
func (self *GuestChangeConfigTask) SetStageComplete(ctx context.Context, data *jsonutils.JSONDict) {
guest := self.GetObject().(*models.SGuest)
hostId := guest.HostId
sessionId, _ := self.Params.GetString("sched_session_id")
lockman.LockRawObject(ctx, models.HostManager.KeywordPlural(), hostId)
defer lockman.ReleaseRawObject(ctx, models.HostManager.KeywordPlural(), hostId)
models.HostManager.ClearSchedDescSessionCache(hostId, sessionId)
self.SSchedTask.SetStageComplete(ctx, data)
}
+29 -14
View File
@@ -104,19 +104,12 @@ func StartScheduleObjects(
doScheduleObjects(ctx, task, schedObjs)
}
func doScheduleObjects(
func doScheduleWithInput(
ctx context.Context,
task IScheduleTask,
objs []IScheduleModel,
) {
schedInput, err := task.GetSchedParams()
if err != nil {
onSchedulerRequestFail(ctx, task, objs, jsonutils.NewString(fmt.Sprintf("GetSchedParams fail: %s", err)))
return
}
//schedInput = models.ApplySchedPolicies(schedInput)
// fetch pendingUsages
schedInput *schedapi.ScheduleInput,
count int,
) (*schedapi.ScheduleOutput, error) {
computeUsage := models.SQuota{}
task.GetPendingUsage(&computeUsage, 0)
regionUsage := models.SRegionQuota{}
@@ -127,11 +120,28 @@ func doScheduleObjects(
jsonutils.Marshal(&regionUsage),
}
params := jsonutils.Marshal(schedInput).(*jsonutils.JSONDict)
var params *jsonutils.JSONDict
if count > 0 {
// if object count <=0, don't need update schedule params
params = jsonutils.Marshal(schedInput).(*jsonutils.JSONDict)
}
task.SetStage("OnScheduleComplete", params)
s := auth.GetSession(ctx, task.GetUserCred(), options.Options.Region, "")
output, err := modules.SchedManager.DoSchedule(s, schedInput, len(objs))
return modules.SchedManager.DoSchedule(s, schedInput, count)
}
func doScheduleObjects(
ctx context.Context,
task IScheduleTask,
objs []IScheduleModel,
) {
schedInput, err := task.GetSchedParams()
if err != nil {
onSchedulerRequestFail(ctx, task, objs, jsonutils.NewString(fmt.Sprintf("GetSchedParams fail: %s", err)))
return
}
output, err := doScheduleWithInput(ctx, task, schedInput, len(objs))
if err != nil {
onSchedulerRequestFail(ctx, task, objs, jsonutils.NewString(err.Error()))
return
@@ -186,6 +196,11 @@ func onSchedulerResults(
objs []IScheduleModel,
results []*schedapi.CandidateResource,
) {
if len(objs) == 0 {
// sched with out object can't clean sched cache immediately
task.SaveScheduleResult(ctx, nil, results[0])
return
}
sort.Sort(sortedIScheduleModelList(objs))
succCount := 0
for idx := 0; idx < len(objs); idx += 1 {