diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 437794b336..25b6059956 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -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, diff --git a/pkg/compute/tasks/guest_change_config_task.go b/pkg/compute/tasks/guest_change_config_task.go index 8bd27af3d1..e2f6779721 100644 --- a/pkg/compute/tasks/guest_change_config_task.go +++ b/pkg/compute/tasks/guest_change_config_task.go @@ -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) +} diff --git a/pkg/compute/tasks/schedule.go b/pkg/compute/tasks/schedule.go index e49d6164ef..222b4567b7 100644 --- a/pkg/compute/tasks/schedule.go +++ b/pkg/compute/tasks/schedule.go @@ -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(®ionUsage), } - 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 {