diff --git a/pkg/apis/compute/guest_const.go b/pkg/apis/compute/guest_const.go index 74c2415bb1..ce373f2358 100644 --- a/pkg/apis/compute/guest_const.go +++ b/pkg/apis/compute/guest_const.go @@ -326,6 +326,8 @@ const ( VM_METADATA_START_VCPU_COUNT = "start_vcpu_count" VM_METADATA_RELEASED_DEVICES = "released_devices" + + VM_METADATA_CPU_NUMA_PIN = "__cpu_numa_pin" ) // windows allow a maximal length of 15 diff --git a/pkg/apis/compute/guests.go b/pkg/apis/compute/guests.go index 6c1f422c01..01fba69e26 100644 --- a/pkg/apis/compute/guests.go +++ b/pkg/apis/compute/guests.go @@ -721,6 +721,7 @@ type ServerMigrateForecastInput struct { SkipKernelCheck bool `json:"skip_kernel_check"` ConvertToKvm bool `json:"convert_to_kvm"` IsRescueMode bool `json:"is_rescue_mode"` + ResetCpuNumaPin bool `json:"reset_cpu_numa_pin"` } type ServerResizeDiskInput struct { @@ -876,6 +877,8 @@ type GuestJsonDesc struct { IsolatedDevices []*IsolatedDeviceJsonDesc `json:"isolated_devices"` + CpuNumaPin []SCpuNumaPin `json:"cpu_numa_pin"` + Domain string `json:"domain"` Nics []*GuestnetworkJsonDesc `json:"nics"` @@ -935,6 +938,18 @@ type GuestJsonDesc struct { Containers []*host.ContainerDesc `json:"containers"` } +type SVCpuPin struct { + Vcpu int + Pcpu int +} + +type SCpuNumaPin struct { + SizeMB *int `json:"size_mb"` + NodeId int `json:"node_id"` + + VcpuPin []SVCpuPin `json:"vcpu_pin"` +} + type ServerSetBootIndexInput struct { // key index, value boot_index Disks map[string]int8 `json:"disks"` diff --git a/pkg/apis/host/types.go b/pkg/apis/host/types.go index d1dadb3083..80ec8270a0 100644 --- a/pkg/apis/host/types.go +++ b/pkg/apis/host/types.go @@ -24,6 +24,11 @@ type ServerCloneDiskFromStorageResponse struct { TargetFormat string `json:"target_format"` } +type HostNodeHugepageNr struct { + NodeId int `json:"node_id"` + HugepageNr int `json:"hugepage_nr"` +} + type HostTopology struct { *topology.Info } diff --git a/pkg/apis/scheduler/api.go b/pkg/apis/scheduler/api.go index 7f8c8164c1..5baf89003d 100644 --- a/pkg/apis/scheduler/api.go +++ b/pkg/apis/scheduler/api.go @@ -78,13 +78,17 @@ type ScheduleInput struct { ServerConfig // HostId used by migrate - HostId string `json:"host_id"` - LiveMigrate bool `json:"live_migrate"` - SkipCpuCheck *bool `json:"skip_cpu_check"` - CpuDesc string `json:"cpu_desc"` - CpuMicrocode string `json:"cpu_microcode"` - CpuMode string `json:"cpu_mode"` - OsArch string `json:"os_arch"` + HostId string `json:"host_id"` + LiveMigrate bool `json:"live_migrate"` + SkipCpuCheck *bool `json:"skip_cpu_check"` + CpuDesc string `json:"cpu_desc"` + CpuMicrocode string `json:"cpu_microcode"` + CpuMode string `json:"cpu_mode"` + OsArch string `json:"os_arch"` + ResetCpuNumaPin bool `json:"reset_cpu_numa_pin"` + + // For Migrate + CpuNumaPin []SCpuNumaPin `json:"cpu_numa_pin"` HostMemPageSizeKB int `json:"host_mem_page_size"` SkipKernelCheck *bool `json:"skip_kernel_check"` @@ -130,12 +134,82 @@ type CandidateNet struct { NetworkIds []string `json:"network_ids"` } +type SFreeNumaCpuMem struct { + FreeCpuCount int + MemSize int + EnableNumaAllocate bool + + CpuCount int + NodeId int +} + +type SortedFreeNumaCpuMam []*SFreeNumaCpuMem + +func (pq SortedFreeNumaCpuMam) Len() int { return len(pq) } + +func (pq SortedFreeNumaCpuMam) Less(i, j int) bool { + if pq[i].EnableNumaAllocate { + return pq[i].MemSize > pq[j].MemSize + } else { + return pq[i].FreeCpuCount > pq[j].FreeCpuCount + } +} + +func (pq SortedFreeNumaCpuMam) Swap(i, j int) { + pq[i], pq[j] = pq[j], pq[i] +} + +func NodesFreeMemSizeEnough(nodeCount, memSize int, cpuNumaFree []*SFreeNumaCpuMem) bool { + if !cpuNumaFree[0].EnableNumaAllocate { + return true + } + + var freeMem = 0 + var leastFree = memSize / nodeCount + for i := 0; i < nodeCount; i++ { + if cpuNumaFree[i].MemSize < leastFree { + return false + } + freeMem += cpuNumaFree[i].MemSize + } + return freeMem >= memSize +} + +func NodesFreeCpuEnough(nodeCount, vcpuCount int, cpuNumaFree []*SFreeNumaCpuMem) bool { + var freeCpu = 0 + var leaseCpu = vcpuCount / nodeCount + + //if vcpuCount > nodeCount*cpuNumaFree[0].CpuCount { + // return false + //} + + for i := 0; i < nodeCount; i++ { + if cpuNumaFree[i].FreeCpuCount < leaseCpu { + return false + } + freeCpu += cpuNumaFree[i].FreeCpuCount + } + return freeCpu >= vcpuCount +} + +type SCpuPin struct { + Vcpu int + Pcpu int +} + +type SCpuNumaPin struct { + CpuPin []int + NodeId int + MemSizeMB *int +} + type CandidateResource struct { - SessionId string `json:"session_id"` - HostId string `json:"host_id"` - Name string `json:"name"` - Disks []*CandidateDisk `json:"disks"` - Nets []*CandidateNet `json:"nets"` + SessionId string `json:"session_id"` + HostId string `json:"host_id"` + Name string `json:"name"` + CpuNumaPin []SCpuNumaPin `json:"cpu_numa_pin"` + Disks []*CandidateDisk `json:"disks"` + Nets []*CandidateNet `json:"nets"` // used by backup schedule BackupCandidate *CandidateResource `json:"backup_candidate"` diff --git a/pkg/cloudcommon/db/opslog_const.go b/pkg/cloudcommon/db/opslog_const.go index cc75f74eda..940386ad31 100644 --- a/pkg/cloudcommon/db/opslog_const.go +++ b/pkg/cloudcommon/db/opslog_const.go @@ -74,6 +74,8 @@ const ( ACT_MIGRATE = "migrate" ACT_MIGRATE_FAIL = "migrate_fail" + ACT_RESET_CPU_NUMA_PIN = "reset_cpu_numa_pin" + ACT_VM_CONVERT = "vm_convert" ACT_VM_CONVERTING = "vm_converting" ACT_VM_CONVERT_FAIL = "vm_convert_fail" diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index 0b1bdd6d14..0921eff935 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -447,7 +447,8 @@ func (self *SKVMGuestDriver) RequestAssociateEip(ctx context.Context, userCred m } func (self *SKVMGuestDriver) RequestChangeVmConfig(ctx context.Context, guest *models.SGuest, task taskman.ITask, instanceType string, vcpuCount, cpuSockets, vmemSize int64) error { - if jsonutils.QueryBoolean(task.GetParams(), "guest_online", false) { + taskParams := task.GetParams() + if jsonutils.QueryBoolean(taskParams, "guest_online", false) { addCpu := vcpuCount - int64(guest.VcpuCount) addMem := vmemSize - int64(guest.VmemSize) if addCpu < 0 || addMem < 0 { @@ -461,6 +462,10 @@ func (self *SKVMGuestDriver) RequestChangeVmConfig(ctx context.Context, guest *m if vmemSize > int64(guest.VmemSize) { body.Set("add_mem", jsonutils.NewInt(addMem)) } + if taskParams.Contains("cpu_numa_pin") { + cpuNumaPin, _ := taskParams.Get("cpu_numa_pin") + body.Set("cpu_numa_pin", cpuNumaPin) + } host, _ := guest.GetHost() url := fmt.Sprintf("%s/servers/%s/hotplug-cpu-mem", host.ManagerUri, guest.Id) _, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, body, false) diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index b51a36d112..f16175a71c 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -527,6 +527,8 @@ func (self *SGuest) GetSchedMigrateParams( if input.PreferHostId != "" { schedDesc.ServerConfig.PreferHost = input.PreferHostId } + + schedDesc.ResetCpuNumaPin = input.ResetCpuNumaPin if input.LiveMigrate { schedDesc.LiveMigrate = input.LiveMigrate if self.GetMetadata(context.Background(), "__cpu_mode", userCred) != api.CPU_MODE_QEMU { @@ -544,6 +546,11 @@ func (self *SGuest) GetSchedMigrateParams( schedDesc.SkipKernelCheck = &input.SkipKernelCheck schedDesc.HostMemPageSizeKB = host.PageSizeKB } + if self.CpuNumaPin != nil { + cpuNumaPin := make([]schedapi.SCpuNumaPin, 0) + self.CpuNumaPin.Unmarshal(&cpuNumaPin) + schedDesc.CpuNumaPin = cpuNumaPin + } } schedDesc.ReuseNetwork = true return schedDesc @@ -562,7 +569,8 @@ func (self *SGuest) StartMigrateTask( ctx context.Context, userCred mcclient.TokenCredential, isRescueMode, autoStart bool, guestStatus, preferHostId, parentTaskId string, ) error { - self.SetStatus(ctx, userCred, api.VM_START_MIGRATE, "") + vmStatus := api.VM_START_MIGRATE + data := jsonutils.NewDict() if isRescueMode { data.Set("is_rescue_mode", jsonutils.JSONTrue) @@ -573,11 +581,17 @@ func (self *SGuest) StartMigrateTask( if autoStart { data.Set("auto_start", jsonutils.JSONTrue) } + if self.HostId == preferHostId { + vmStatus = api.VM_STARTING + data.Set("reset_cpu_numa_pin", jsonutils.JSONTrue) + } + data.Set("guest_status", jsonutils.NewString(guestStatus)) dedicateMigrateTask := "GuestMigrateTask" if self.GetHypervisor() != api.HYPERVISOR_KVM { dedicateMigrateTask = "ManagedGuestMigrateTask" //托管私有云 } + self.SetStatus(ctx, userCred, vmStatus, "") if task, err := taskman.TaskManager.NewTask(ctx, dedicateMigrateTask, self, userCred, data, parentTaskId, "", nil); err != nil { log.Errorln(err) return err @@ -1543,7 +1557,12 @@ func (self *SGuest) StartGueststartTask( ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict, parentTaskId string, ) error { - if self.Hypervisor == api.HYPERVISOR_KVM && self.guestDisksStorageTypeIsShared() { + schedStart := self.Hypervisor == api.HYPERVISOR_KVM && self.guestDisksStorageTypeIsShared() + if options.Options.IgnoreNonrunningGuests && self.CpuNumaPin != nil { + schedStart = true + } + + if schedStart { return self.GuestSchedStartTask(ctx, userCred, data, parentTaskId) } else { return self.GuestNonSchedStartTask(ctx, userCred, data, parentTaskId) diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 8f285dac38..ff03bd857b 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -127,6 +127,8 @@ type SGuest struct { VcpuCount int `nullable:"false" default:"1" list:"user" create:"optional"` // 内存大小, 单位MB VmemSize int `nullable:"false" list:"user" create:"required"` + // CPU 内存绑定信息 + CpuNumaPin jsonutils.JSONObject `nullable:"true" get:"user" update:"user" create:"optional"` // 启动顺序 BootOrder string `width:"8" charset:"ascii" nullable:"true" default:"cdn" list:"user" update:"user" create:"optional"` @@ -1271,6 +1273,94 @@ func (guest *SGuest) SetHostIdWithBackup(userCred mcclient.TokenCredential, mast return err } +func (guest *SGuest) UpdateCpuNumaPin( + ctx context.Context, userCred mcclient.TokenCredential, + schedCpuNumaPin []schedapi.SCpuNumaPin, cpuNumaPinTarget []api.SCpuNumaPin, +) error { + srcSchedCpuNumaPin := make([]schedapi.SCpuNumaPin, 0) + err := guest.CpuNumaPin.Unmarshal(&srcSchedCpuNumaPin) + if err != nil { + return err + } + srcSchedCpuNumaPin = append(srcSchedCpuNumaPin, schedCpuNumaPin...) + + srcCpuNumaPin := make([]api.SCpuNumaPin, 0) + cpuNumaPinStr := guest.GetMetadata(ctx, api.VM_METADATA_CPU_NUMA_PIN, nil) + cpuNumaPinJson, err := jsonutils.ParseString(cpuNumaPinStr) + if err != nil { + return err + } + err = cpuNumaPinJson.Unmarshal(&srcCpuNumaPin) + if err != nil { + return err + } + srcCpuNumaPin = append(srcCpuNumaPin, cpuNumaPinTarget...) + + diff, err := db.Update(guest, func() error { + guest.CpuNumaPin = jsonutils.Marshal(srcSchedCpuNumaPin) + return nil + }) + if err != nil { + return err + } + + var jcpuNumaPin = jsonutils.Marshal(srcCpuNumaPin) + err = guest.SetMetadata(ctx, api.VM_METADATA_CPU_NUMA_PIN, jcpuNumaPin, userCred) + if err != nil { + return err + } + + db.OpsLog.LogEvent(guest, db.ACT_UPDATE, diff, userCred) + return nil +} + +func (guest *SGuest) SetCpuNumaPin( + ctx context.Context, userCred mcclient.TokenCredential, + schedCpuNumaPin []schedapi.SCpuNumaPin, cpuNumaPin []api.SCpuNumaPin, +) error { + if cpuNumaPin == nil && schedCpuNumaPin != nil { + cpuNumaPin = make([]api.SCpuNumaPin, len(schedCpuNumaPin)) + + vcpuId := 0 + for i := range schedCpuNumaPin { + cpuNumaPin[i] = api.SCpuNumaPin{ + SizeMB: schedCpuNumaPin[i].MemSizeMB, + NodeId: schedCpuNumaPin[i].NodeId, + } + cpuNumaPin[i].VcpuPin = make([]api.SVCpuPin, len(schedCpuNumaPin[i].CpuPin)) + for j := range schedCpuNumaPin[i].CpuPin { + cpuNumaPin[i].VcpuPin[j].Pcpu = schedCpuNumaPin[i].CpuPin[j] + cpuNumaPin[i].VcpuPin[j].Vcpu = vcpuId + vcpuId += 1 + } + } + } + + var schedCpuNumaPinJ jsonutils.JSONObject + if schedCpuNumaPin != nil { + schedCpuNumaPinJ = jsonutils.Marshal(schedCpuNumaPin) + } + diff, err := db.Update(guest, func() error { + guest.CpuNumaPin = schedCpuNumaPinJ + return nil + }) + if err != nil { + return err + } + + var jcpuNumaPin interface{} = "" + if cpuNumaPin != nil { + jcpuNumaPin = jsonutils.Marshal(cpuNumaPin) + } + err = guest.SetMetadata(ctx, api.VM_METADATA_CPU_NUMA_PIN, jcpuNumaPin, userCred) + if err != nil { + return err + } + + db.OpsLog.LogEvent(guest, db.ACT_UPDATE, diff, userCred) + return err +} + func (guest *SGuest) ValidateResizeDisk(disk *SDisk, storage *SStorage) error { drv, err := guest.GetDriver() if err != nil { @@ -5105,6 +5195,18 @@ func (self *SGuest) GetJsonDescAtHypervisor(ctx context.Context, host *SHost) *a desc.IsolatedDevices = append(desc.IsolatedDevices, dev.getDesc()) } + if self.CpuNumaPin != nil { + cpuNumaPin := make([]api.SCpuNumaPin, 0) + cpuNumaPinStr := self.GetMetadata(ctx, api.VM_METADATA_CPU_NUMA_PIN, nil) + cpuNumaPinJson, err := jsonutils.ParseString(cpuNumaPinStr) + if err != nil { + log.Errorf("failed parse cpu numa pin %s: %s", cpuNumaPinStr, err) + } else { + cpuNumaPinJson.Unmarshal(&cpuNumaPin) + desc.CpuNumaPin = cpuNumaPin + } + } + // nics, domain desc.Domain = options.Options.DNSDomain nics, _ := self.GetNetworks("") diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index ea95e59211..b44b32d049 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -4250,7 +4250,7 @@ func (hh *SHost) ValidateUpdateData(ctx context.Context, userCred mcclient.Token func (hh *SHost) PostUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) { hh.SEnabledStatusInfrasResourceBase.PostUpdate(ctx, userCred, query, data) - if data.Contains("cpu_cmtbound") || data.Contains("mem_cmtbound") { + if data.Contains("cpu_cmtbound") || data.Contains("mem_cmtbound") || data.Contains("enable_numa_allocate") { hh.ClearSchedDescCache() } diff --git a/pkg/compute/tasks/guest_batch_create_task.go b/pkg/compute/tasks/guest_batch_create_task.go index 88027eb721..2a33af6d05 100644 --- a/pkg/compute/tasks/guest_batch_create_task.go +++ b/pkg/compute/tasks/guest_batch_create_task.go @@ -138,6 +138,14 @@ func (task *GuestBatchCreateTask) allocateGuestOnHost(ctx context.Context, guest } } + if len(candidate.CpuNumaPin) > 0 { + if err := guest.SetCpuNumaPin(ctx, task.UserCred, candidate.CpuNumaPin, nil); err != nil { + log.Errorf("SetCpuNumaPin fail %s", err) + guest.SetStatus(ctx, task.UserCred, api.VM_CREATE_FAILED, err.Error()) + return err + } + } + host, _ := guest.GetHost() quotaCpuMem := models.SQuota{Count: 1, Cpu: int(guest.VcpuCount), Memory: guest.VmemSize} diff --git a/pkg/compute/tasks/guest_change_config_task.go b/pkg/compute/tasks/guest_change_config_task.go index a1617067f1..3262e235f4 100644 --- a/pkg/compute/tasks/guest_change_config_task.go +++ b/pkg/compute/tasks/guest_change_config_task.go @@ -85,6 +85,11 @@ func (task *GuestChangeConfigTask) SaveScheduleResult(ctx context.Context, obj I // must get object from task, because of obj is nil guest := task.GetObject().(*models.SGuest) task.Params.Set("sched_session_id", jsonutils.NewString(target.SessionId)) + + if len(target.CpuNumaPin) > 0 { + task.Params.Set("cpu_numa_pin", jsonutils.Marshal(target.CpuNumaPin)) + } + /*confs, err := task.getChangeConfigSetting() if err != nil { task.markStageFailed(ctx, guest, jsonutils.NewString(err.Error())) @@ -258,6 +263,19 @@ func (task *GuestChangeConfigTask) OnGuestChangeCpuMemSpecComplete(ctx context.C task.markStageFailed(ctx, guest, jsonutils.NewString(fmt.Sprintf("Update fail %s", err))) return } + + if task.Params.Contains("cpu_numa_pin") { + cpuNumaPinSched := make([]schedapi.SCpuNumaPin, 0) + task.Params.Unmarshal(&cpuNumaPinSched, "cpu_numa_pin") + cpuNumaPinTarget := make([]api.SCpuNumaPin, 0) + data.Unmarshal(&cpuNumaPinTarget, "cpu_numa_pin") + err = guest.UpdateCpuNumaPin(ctx, task.UserCred, cpuNumaPinSched, cpuNumaPinTarget) + if err != nil { + task.markStageFailed(ctx, guest, jsonutils.NewString(fmt.Sprintf("Update cpu numa pin fail %s", err))) + return + } + } + changeConfigSpec := guest.GetShortDesc(ctx) if confs.VcpuCount > 0 && confs.AddedCpu() > 0 { changeConfigSpec.Set("add_cpu", jsonutils.NewInt(int64(confs.AddedCpu()))) diff --git a/pkg/compute/tasks/guest_live_migrate_task.go b/pkg/compute/tasks/guest_live_migrate_task.go index 343189f2be..f7db6dabe7 100644 --- a/pkg/compute/tasks/guest_live_migrate_task.go +++ b/pkg/compute/tasks/guest_live_migrate_task.go @@ -74,6 +74,10 @@ func (task *GuestMigrateTask) GetSchedParams() (*schedapi.ScheduleInput, error) preferHostId, _ := task.Params.GetString("prefer_host_id") input.PreferHostId = preferHostId } + if jsonutils.QueryBoolean(task.Params, "reset_cpu_numa_pin", false) { + input.ResetCpuNumaPin = true + } + if task.isLiveMigrate() { input.LiveMigrate = true skipCpuCheck := jsonutils.QueryBoolean(task.Params, "skip_cpu_check", false) @@ -105,19 +109,29 @@ func (task *GuestMigrateTask) OnScheduleFailed(ctx context.Context, reason jsonu } func (task *GuestMigrateTask) SaveScheduleResult(ctx context.Context, obj IScheduleModel, target *schedapi.CandidateResource, index int) { - targetHostId := target.HostId guest := obj.(*models.SGuest) + if jsonutils.QueryBoolean(task.Params, "reset_cpu_numa_pin", false) { + guest.SetCpuNumaPin(ctx, task.UserCred, target.CpuNumaPin, nil) + db.OpsLog.LogEvent(guest, db.ACT_RESET_CPU_NUMA_PIN, fmt.Sprintf("reset cpu numa pin %s", jsonutils.Marshal(target.CpuNumaPin)), task.UserCred) + task.SetStageComplete(ctx, nil) + return + } + + targetHostId := target.HostId targetHost := models.HostManager.FetchHostById(targetHostId) if targetHost == nil { task.TaskFailed(ctx, guest, jsonutils.NewString("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), task.UserCred) - logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_MIGRATING, - fmt.Sprintf("guest start migrate from host %s to %s(%s)", guest.HostId, targetHostId, targetHost.GetName()), task.UserCred, true) body := jsonutils.NewDict() body.Set("target_host_id", jsonutils.NewString(targetHostId)) + if len(target.CpuNumaPin) > 0 { + body.Set("target_cpu_numa_pin", jsonutils.Marshal(target.CpuNumaPin)) + } else { + body.Set("target_cpu_numa_pin", jsonutils.JSONNull) + } + // for params notes body.Set("target_host_name", jsonutils.NewString(targetHost.Name)) srcHost := models.HostManager.FetchHostById(guest.HostId) @@ -166,6 +180,11 @@ func (task *GuestMigrateTask) SaveScheduleResult(ctx context.Context, obj ISched body.Set("cache_templates", jsonutils.NewStringArray(templates)) } } + + db.OpsLog.LogEvent(guest, db.ACT_MIGRATING, fmt.Sprintf("guest start migrate from host %s to %s", guest.HostId, targetHostId), task.UserCred) + logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_MIGRATING, + fmt.Sprintf("guest start migrate from host %s to %s(%s)", guest.HostId, targetHostId, targetHost.GetName()), task.UserCred, true) + task.SetStage("OnStartCacheImages", body) task.OnStartCacheImages(ctx, guest, nil) } @@ -434,6 +453,12 @@ func (task *GuestMigrateTask) sharedStorageMigrateConf(ctx context.Context, gues body.Set("is_local_storage", jsonutils.JSONFalse) body.Set("qemu_version", jsonutils.NewString(guest.GetQemuVersion(task.UserCred))) targetDesc := guest.GetJsonDescAtHypervisor(ctx, targetHost) + if task.Params.Contains("target_cpu_numa_pin") { + if err := task.setCpuNumaPin(targetDesc); err != nil { + return nil, errors.Wrap(err, "setCpuNumaPin") + } + } + body.Set("desc", jsonutils.Marshal(targetDesc)) sourceHost, _ := guest.GetHost() @@ -488,6 +513,12 @@ func (task *GuestMigrateTask) localStorageMigrateConf(ctx context.Context, if len(targetDesc.Disks) == 0 { return nil, errors.Errorf("Get disksDesc error") } + if task.Params.Contains("target_cpu_numa_pin") { + if err := task.setCpuNumaPin(targetDesc); err != nil { + return nil, errors.Wrap(err, "setCpuNumaPin") + } + } + targetStorages, _ := task.Params.GetArray("target_storages") for i := 0; i < len(disks); i++ { targetStorageId, err := targetStorages[i].GetString() @@ -503,6 +534,20 @@ func (task *GuestMigrateTask) localStorageMigrateConf(ctx context.Context, return body, nil } +func (task *GuestMigrateTask) setCpuNumaPin(targetDesc *api.GuestJsonDesc) error { + cpuNumaPin := make([]schedapi.SCpuNumaPin, 0) + if err := task.Params.Unmarshal(&cpuNumaPin, "cpu_numa_pin"); err != nil { + return errors.Wrap(err, "unmarshal cpu_numa_pin") + } + for i := range targetDesc.CpuNumaPin { + for j := range targetDesc.CpuNumaPin[i].VcpuPin { + targetDesc.CpuNumaPin[i].VcpuPin[j].Pcpu = cpuNumaPin[i].CpuPin[j] + } + } + task.Params.Set("target_vcpu_numa_pin", jsonutils.Marshal(targetDesc.CpuNumaPin)) + return nil +} + func (task *GuestLiveMigrateTask) OnStartDestComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { liveMigrateDestPort, err := data.Get("live_migrate_dest_port") if err != nil { @@ -575,6 +620,29 @@ func (task *GuestMigrateTask) setGuest(ctx context.Context, guest *models.SGuest } } } + + if task.Params.Contains("target_cpu_numa_pin") { + var cpuNumaPinSrc []schedapi.SCpuNumaPin = nil + var cpuNumaPin []api.SCpuNumaPin = nil + + val, _ := task.Params.Get("target_cpu_numa_pin") + if !val.Equals(jsonutils.JSONNull) { + cpuNumaPinSrc = make([]schedapi.SCpuNumaPin, 0) + if err := task.Params.Unmarshal(&cpuNumaPinSrc, "target_cpu_numa_pin"); err != nil { + return errors.Wrap(err, "unmarshal target_cpu_numa_pin") + } + + cpuNumaPin = make([]api.SCpuNumaPin, 0) + if err := task.Params.Unmarshal(&cpuNumaPin, "target_vcpu_numa_pin"); err != nil { + return errors.Wrap(err, "unmarshal target_vcpu_numa_pin") + } + } + + if err := guest.SetCpuNumaPin(ctx, task.UserCred, cpuNumaPinSrc, cpuNumaPin); err != nil { + return errors.Wrap(err, "SetCpuNumaPin") + } + } + oldHost, _ := guest.GetHost() oldHost.ClearSchedDescCache() err := guest.OnScheduleToHost(ctx, task.UserCred, targetHostId) diff --git a/pkg/compute/tasks/guest_start_task.go b/pkg/compute/tasks/guest_start_task.go index 7e06a19271..64cdd16eab 100644 --- a/pkg/compute/tasks/guest_start_task.go +++ b/pkg/compute/tasks/guest_start_task.go @@ -152,6 +152,11 @@ func (self *GuestSchedStartTask) OnInit(ctx context.Context, obj db.IStandaloneM } func (self *GuestSchedStartTask) StartScheduler(ctx context.Context, guest *models.SGuest) { + if guest.CpuNumaPin != nil { + self.ScheduleFailed(ctx, guest) + return + } + host, _ := guest.GetHost() if request := host.GetRunningGuestResourceUsage(); request == nil { self.TaskFailed(ctx, guest, jsonutils.NewString("Guest Start Failed: Can't Get Host Guests CPU Memory Usage")) @@ -170,7 +175,13 @@ func (self *GuestSchedStartTask) StartScheduler(ctx context.Context, guest *mode func (self *GuestSchedStartTask) ScheduleFailed(ctx context.Context, guest *models.SGuest) { self.SetStage("OnGuestMigrate", nil) - guest.StartMigrateTask(ctx, self.UserCred, false, false, guest.Status, "", self.GetId()) + + preferHostId := "" + if guest.CpuNumaPin != nil { + preferHostId = guest.HostId + } + + guest.StartMigrateTask(ctx, self.UserCred, false, false, guest.Status, preferHostId, self.GetId()) } func (self *GuestSchedStartTask) OnGuestMigrate(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { diff --git a/pkg/hostman/guestman/desc/desc.go b/pkg/hostman/guestman/desc/desc.go index 031ced4c37..7e8ca55cdc 100644 --- a/pkg/hostman/guestman/desc/desc.go +++ b/pkg/hostman/guestman/desc/desc.go @@ -57,10 +57,15 @@ type SMemSlot struct { type SCpuNumaPin struct { SizeMB int64 - Regular bool - HostNodes *uint16 `json:",omitempty"` - Vcpus *string `json:",omitempty"` - Pcpus *string `json:",omitempty"` + Unregular bool + NodeId *uint16 `json:",omitempty"` + + VcpuPin []SVCpuPin `json:",omitempty"` +} + +type SVCpuPin struct { + Vcpu int + Pcpu int } type SMemDesc struct { diff --git a/pkg/hostman/guestman/guesthandlers/guesthandler.go b/pkg/hostman/guestman/guesthandlers/guesthandler.go index f3d7e0a342..50c8f0525c 100644 --- a/pkg/hostman/guestman/guesthandlers/guesthandler.go +++ b/pkg/hostman/guestman/guesthandlers/guesthandler.go @@ -25,6 +25,7 @@ import ( computeapi "yunion.io/x/onecloud/pkg/apis/compute" hostapi "yunion.io/x/onecloud/pkg/apis/host" + schedapi "yunion.io/x/onecloud/pkg/apis/scheduler" "yunion.io/x/onecloud/pkg/appsrv" "yunion.io/x/onecloud/pkg/hostman/guestman" "yunion.io/x/onecloud/pkg/hostman/guestman/desc" @@ -581,12 +582,39 @@ func guestHotplugCpuMem(ctx context.Context, userCred mcclient.TokenCredential, addCpuCount, _ := body.Int("add_cpu") addMemSize, _ := body.Int("add_mem") - hostutils.DelayTaskWithoutReqctx(ctx, guestman.GetGuestManager().HotplugCpuMem, - &guestman.SGuestHotplugCpuMem{ - Sid: sid, - AddCpuCount: addCpuCount, - AddMemSize: addMemSize, - }) + + input := &guestman.SGuestHotplugCpuMem{ + Sid: sid, + AddCpuCount: addCpuCount, + AddMemSize: addMemSize, + } + + if body.Contains("cpu_numa_pin") { + cpuNumaPin := make([]schedapi.SCpuNumaPin, 0) + if err := body.Unmarshal(&cpuNumaPin, "cpu_numa_pin"); err != nil { + return nil, httperrors.NewInputParameterError("failed parse cpu_numa_pin %s", err) + } + + if len(cpuNumaPin) > 0 { + descCpuNumaPin := make([]*desc.SCpuNumaPin, len(cpuNumaPin)) + for i := range cpuNumaPin { + if cpuNumaPin[i].MemSizeMB != nil { + descCpuNumaPin[i].SizeMB = int64(*cpuNumaPin[i].MemSizeMB) + nodeId := uint16(cpuNumaPin[i].NodeId) + descCpuNumaPin[i].NodeId = &nodeId + } + if len(cpuNumaPin[i].CpuPin) > 0 { + vcpuPin := make([]desc.SVCpuPin, len(cpuNumaPin[i].CpuPin)) + for j := range vcpuPin { + vcpuPin[j].Pcpu = cpuNumaPin[i].CpuPin[j] + } + descCpuNumaPin[i].VcpuPin = vcpuPin + } + } + input.CpuNumaPin = descCpuNumaPin + } + } + hostutils.DelayTaskWithoutReqctx(ctx, guestman.GetGuestManager().HotplugCpuMem, input) return nil, nil } diff --git a/pkg/hostman/guestman/guesthelper.go b/pkg/hostman/guestman/guesthelper.go index 40acfe745e..44506f9868 100644 --- a/pkg/hostman/guestman/guesthelper.go +++ b/pkg/hostman/guestman/guesthelper.go @@ -100,6 +100,8 @@ type SGuestHotplugCpuMem struct { Sid string AddCpuCount int64 AddMemSize int64 + + CpuNumaPin []*desc.SCpuNumaPin } type SReloadDisk struct { @@ -214,7 +216,7 @@ type CpuSetCounter struct { Lock sync.Mutex } -func NewGuestCpuSetCounter(info *hostapi.HostTopology, reservedCpus *cpuset.CPUSet, numaAllocate bool, hugepageSizeKB int) (*CpuSetCounter, error) { +func NewGuestCpuSetCounter(info *hostapi.HostTopology, reservedCpus *cpuset.CPUSet, numaAllocate bool, hugepageSizeKB, cpuCmtbound int) (*CpuSetCounter, error) { cpuSetCounter := new(CpuSetCounter) cpuSetCounter.Nodes = make([]*NumaNode, len(info.Nodes)) cpuSetCounter.NumaEnabled = numaAllocate @@ -241,6 +243,8 @@ func NewGuestCpuSetCounter(info *hostapi.HostTopology, reservedCpus *cpuset.CPUS cpuDie.LogicalProcessors = dieBuilder.Result() node.CpuCount += cpuDie.LogicalProcessors.Size() node.LogicalProcessors = node.LogicalProcessors.Union(cpuDie.LogicalProcessors) + cpuDie.initCpuFree(cpuCmtbound) + cpuDies = append(cpuDies, cpuDie) } if !hasL3Cache { @@ -257,11 +261,14 @@ func NewGuestCpuSetCounter(info *hostapi.HostTopology, reservedCpus *cpuset.CPUS cpuDie.LogicalProcessors = dieBuilder.Result() node.CpuCount += cpuDie.LogicalProcessors.Size() node.LogicalProcessors = node.LogicalProcessors.Union(cpuDie.LogicalProcessors) + cpuDie.initCpuFree(cpuCmtbound) + cpuDies = append(cpuDies, cpuDie) } hasL3Cache = false node.CpuDies = cpuDies + sort.Sort(node.CpuDies) cpuSetCounter.Nodes[i] = node } sort.Sort(cpuSetCounter) @@ -292,9 +299,9 @@ func (pq *CpuSetCounter) AllocCpusetWithNodeCount(vcpuCount int, memSizeKB int64 remPcpuCount -= 1 } res[pq.Nodes[i].NodeId] = SAllocNumaCpus{ - Cpuset: pq.Nodes[i].AllocCpuset(npcpuCount), + Cpuset: pq.Nodes[i].AllocCpuset1(npcpuCount), MemSizeKB: nodeAllocSize, - Regular: true, + Unregular: false, } pq.Nodes[i].NumaHugeFreeMemSizeKB -= nodeAllocSize pq.Nodes[i].VcpuCount += npcpuCount @@ -307,7 +314,7 @@ type SAllocNumaCpus struct { Cpuset []int MemSizeKB int64 - Regular bool + Unregular bool } func (pq *CpuSetCounter) IsNumaEnabled() bool { @@ -351,9 +358,9 @@ func (pq *CpuSetCounter) AllocNumaNodes(vcpuCount int, memSizeKB int64, perferNu } if pq.Nodes[i].NumaHugeFreeMemSizeKB >= memSizeKB { res[pq.Nodes[i].NodeId] = SAllocNumaCpus{ - Cpuset: pq.Nodes[i].AllocCpuset(vcpuCount), + Cpuset: pq.Nodes[i].AllocCpuset1(vcpuCount), MemSizeKB: memSizeKB, - Regular: true, + Unregular: false, } pq.Nodes[i].NumaHugeFreeMemSizeKB -= memSizeKB pq.Nodes[i].VcpuCount += vcpuCount @@ -388,9 +395,9 @@ func (pq *CpuSetCounter) AllocNumaNodes(vcpuCount int, memSizeKB int64, perferNu remPcpuCount -= 1 } res[pq.Nodes[i].NodeId] = SAllocNumaCpus{ - Cpuset: pq.Nodes[i].AllocCpuset(npcpuCount), + Cpuset: pq.Nodes[i].AllocCpuset1(npcpuCount), MemSizeKB: nodeAllocSize, - Regular: true, + Unregular: false, } pq.Nodes[i].NumaHugeFreeMemSizeKB -= nodeAllocSize pq.Nodes[i].VcpuCount += npcpuCount @@ -431,9 +438,9 @@ func (pq *CpuSetCounter) setNumaNodes(numaMaps map[int]int, vcpuCount int64) map allocMem := int64(size) * 1024 //npcpuCount := int(vcpuCount*allocMem/memSizeKB + (vcpuCount*allocMem)%memSizeKB) res[pq.Nodes[i].NodeId] = SAllocNumaCpus{ - Cpuset: pq.Nodes[i].AllocCpuset(int(vcpuCount)), + Cpuset: pq.Nodes[i].AllocCpuset1(int(vcpuCount)), MemSizeKB: allocMem, - Regular: false, + Unregular: true, } pq.Nodes[i].NumaHugeFreeMemSizeKB -= allocMem @@ -583,6 +590,62 @@ func NewNumaNode(nodeId int, numaAllocate bool, hugepageSizeKB int) (*NumaNode, return n, nil } +func (n *NumaNode) AllocCpuset1(vcpuCount int) []int { + var allocCount = vcpuCount + var dieCnt = 0 + + // If request vcpu count great then node cpucount, + // vcpus should evenly distributed to all dies. + // Otherwise figure out how many dies can hold + // all of vcpus at first, and evenly distributed + // to selected dies. + if vcpuCount > n.CpuCount { + dieCnt = len(n.CpuDies) + } else { + var pcpuCount = 0 + for dieCnt < len(n.CpuDies) { + pcpuCount += n.CpuDies[dieCnt].LogicalProcessors.Size() + dieCnt += 1 + + if pcpuCount >= vcpuCount { + break + } + } + } + + var perDieCpuCount = vcpuCount / dieCnt + var allocCpuCountMap = make([]int, dieCnt) + for allocCount > 0 { + for i := 0; i < dieCnt; i++ { + var allocNum = perDieCpuCount + if allocCount < allocNum { + allocNum = allocCount + } + allocCount -= allocNum + allocCpuCountMap[i] += allocNum + } + } + + var pcpus = make([]int, 0) + for i := 0; i < len(allocCpuCountMap); i++ { + var allocCpuCount = allocCpuCountMap[i] + for allocCpuCount > 0 { + pcpus := n.CpuDies[i].LogicalProcessors.ToSliceNoSort() + for j := 0; j < len(pcpus); j++ { + if n.CpuDies[i].CpuFree[pcpus[j]] > 0 { + pcpus = append(pcpus, n.CpuDies[i].CpuFree[pcpus[j]]) + n.CpuDies[i].CpuFree[pcpus[j]] -= 1 + } + allocCpuCount -= 1 + if allocCpuCount <= 0 { + break + } + } + } + } + return pcpus +} + func (n *NumaNode) AllocCpuset(vcpuCount int) []int { cpus := make([]int, 0) @@ -602,10 +665,19 @@ func (n *NumaNode) AllocCpuset(vcpuCount int) []int { } type CPUDie struct { + CpuFree map[int]int LogicalProcessors cpuset.CPUSet VcpuCount int } +func (d *CPUDie) initCpuFree(cpuCmtbound int) { + cpuFree := map[int]int{} + for _, cpuId := range d.LogicalProcessors.ToSliceNoSort() { + cpuFree[cpuId] = cpuCmtbound + } + d.CpuFree = cpuFree +} + type SorttedCPUDie []*CPUDie func (pq SorttedCPUDie) Len() int { return len(pq) } @@ -648,7 +720,11 @@ func (pq *SorttedCPUDie) ReleaseCpus(cpus []int, vcpuCount int) { for i := 0; i < len(*pq); i++ { if _, ok := cpuDies[i]; ok { - (*pq)[i].VcpuCount -= vcpuCount + d := (*pq)[i] + for _, cpu := range cpus { + d.CpuFree[cpu] += 1 + } + d.VcpuCount -= vcpuCount } } sort.Sort(pq) @@ -670,8 +746,12 @@ func (pq *SorttedCPUDie) LoadCpus(cpus []int, vcpuCount int) { } for i := 0; i < len(*pq); i++ { - if _, ok := cpuDies[i]; ok { - (*pq)[i].VcpuCount += vcpuCount + if cpus, ok := cpuDies[i]; ok { + d := (*pq)[i] + for _, cpu := range cpus { + d.CpuFree[cpu] -= 1 + } + d.VcpuCount += vcpuCount } } sort.Sort(pq) diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index efd3b190d1..abf63b545b 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -255,9 +255,13 @@ func (m *SGuestManager) CleanServer(sid string) { func (m *SGuestManager) Bootstrap() (chan struct{}, error) { hostTypo := m.host.GetHostTopology() - m.numaAllocate = m.host.IsNumaAllocateEnabled() && m.host.IsHugepagesEnabled() && (len(hostTypo.Nodes) > 1) + + if options.HostOptions.EnableHostAgentNumaAllocate { + m.numaAllocate = m.host.IsNumaAllocateEnabled() && m.host.IsHugepagesEnabled() && (len(hostTypo.Nodes) > 1) + } + cpuSet, err := NewGuestCpuSetCounter( - hostTypo, m.host.GetReservedCpusInfo(), m.numaAllocate, m.host.HugepageSizeKb()) + hostTypo, m.host.GetReservedCpusInfo(), m.numaAllocate, m.host.HugepageSizeKb(), m.host.CpuCmtBound()) if err != nil { return nil, err } @@ -1445,7 +1449,7 @@ func (m *SGuestManager) HotplugCpuMem(ctx context.Context, params interface{}) ( return nil, hostutils.ParamsError } guest, _ := m.GetKVMServer(hotplugParams.Sid) - NewGuestHotplugCpuMemTask(ctx, guest, int(hotplugParams.AddCpuCount), int(hotplugParams.AddMemSize)).Start() + NewGuestHotplugCpuMemTask(ctx, guest, int(hotplugParams.AddCpuCount), int(hotplugParams.AddMemSize), hotplugParams.CpuNumaPin).Start() return nil, nil } diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index 350ebd06b1..47e59d6fcc 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -2393,27 +2393,31 @@ func (task *SGuestOnlineResizeDiskTask) OnResizeSucc(err string) { type SGuestHotplugCpuMemTask struct { *SKVMGuestInstance - ctx context.Context - addCpuCount int - addMemSize int + ctx context.Context + addCpuCount int + addMemSize int + addMemNodeIndex int originalCpuCount int addedCpuCount int addedVcpuIds []int + cpuNumaPin []*desc.SCpuNumaPin - addedMemSize int - memSlotNewIndex *int - memSlot *desc.SMemSlot + addedMemSize int + memSlotNewIndex *int + memSlotNewIndexs []int + memSlots []*desc.SMemSlot } func NewGuestHotplugCpuMemTask( - ctx context.Context, s *SKVMGuestInstance, addCpuCount, addMemSize int, + ctx context.Context, s *SKVMGuestInstance, addCpuCount, addMemSize int, cpuNumaPin []*desc.SCpuNumaPin, ) *SGuestHotplugCpuMemTask { return &SGuestHotplugCpuMemTask{ SKVMGuestInstance: s, ctx: ctx, addCpuCount: addCpuCount, addMemSize: addMemSize, + cpuNumaPin: cpuNumaPin, } } @@ -2446,12 +2450,14 @@ func (task *SGuestHotplugCpuMemTask) buildVcpusMap() { } } + allocatedVcpus := make([]int, 0) for i := range task.Desc.CpuNumaPin { - if task.Desc.CpuNumaPin[i].Vcpus != nil { - allocedVcpus, _ := cpuset.Parse(*task.Desc.CpuNumaPin[i].Vcpus) - vcpuSet = vcpuSet.Difference(allocedVcpus) + for j := range task.Desc.CpuNumaPin[i].VcpuPin { + allocatedVcpus = append(allocatedVcpus, task.Desc.CpuNumaPin[i].VcpuPin[j].Vcpu) } } + allocatedCpuset := cpuset.NewCPUSet(allocatedVcpus...) + vcpuSet = vcpuSet.Difference(allocatedCpuset) task.startAddCpusWithFreeVcpuSet(vcpuSet.ToSlice()) } @@ -2470,17 +2476,29 @@ func (task *SGuestHotplugCpuMemTask) startAddCpusWithFreeVcpuSet(vcpuSet []int) return } - cpus, _ := task.manager.cpuSet.AllocCpuset(1, 0, -1) - for _, cpus := range cpus { - pcpus := cpuset.NewCPUSet(cpus.Cpuset...).String() - vcpus := fmt.Sprintf("%d-%d", vcpuId, vcpuId) - cpuPin := &desc.SCpuNumaPin{ - SizeMB: 0, - Pcpus: &pcpus, - Vcpus: &vcpus, - Regular: true, + if len(task.cpuNumaPin) > 0 { + for i := range task.cpuNumaPin { + for j := range task.cpuNumaPin[i].VcpuPin { + if task.cpuNumaPin[i].VcpuPin[j].Vcpu == -1 { + task.cpuNumaPin[i].VcpuPin[j].Vcpu = vcpuId + } + } + } + } else { + cpus, _ := task.manager.cpuSet.AllocCpuset(1, 0, -1) + for _, cpus := range cpus { + //pcpus := cpuset.NewCPUSet(cpus.Cpuset...).String() + //vcpus := fmt.Sprintf("%d-%d", vcpuId, vcpuId) + vcpuPin := make([]desc.SVCpuPin, 1) + vcpuPin[0].Pcpu = cpus.Cpuset[0] + vcpuPin[0].Vcpu = vcpuId + cpuPin := &desc.SCpuNumaPin{ + SizeMB: 0, + VcpuPin: vcpuPin, + Unregular: false, + } + task.Desc.CpuNumaPin = append(task.Desc.CpuNumaPin, cpuPin) } - task.Desc.CpuNumaPin = append(task.Desc.CpuNumaPin, cpuPin) } if task.addedVcpuIds == nil { task.addedVcpuIds = []int{vcpuId} @@ -2529,6 +2547,21 @@ func (task *SGuestHotplugCpuMemTask) onGetSlotIndex(index int) { var newIndex = index + len(task.Desc.MemDesc.Mem.Mems) task.memSlotNewIndex = &newIndex + var addMemSize = task.addMemSize + var numaNodeDesc *desc.SCpuNumaPin + for i := task.addMemNodeIndex; i < len(task.cpuNumaPin); i++ { + if task.cpuNumaPin[i].SizeMB > 0 { + task.addMemNodeIndex = i + 1 + numaNodeDesc = task.cpuNumaPin[i] + addMemSize = int(task.cpuNumaPin[i].SizeMB) + break + } + } + var hostNodes *uint16 + if numaNodeDesc != nil { + hostNodes = numaNodeDesc.NodeId + } + var objType string var id = fmt.Sprintf("mem%d", *task.memSlotNewIndex) var opts map[string]string @@ -2544,7 +2577,7 @@ func (task *SGuestHotplugCpuMemTask) onGetSlotIndex(index int) { return } err = procutils.NewRemoteCommandAsFarAsPossible("mount", "-t", "hugetlbfs", "-o", - fmt.Sprintf("pagesize=%dK,size=%dM", task.manager.host.HugepageSizeKb(), task.addMemSize), + fmt.Sprintf("pagesize=%dK,size=%dM", task.manager.host.HugepageSizeKb(), addMemSize), fmt.Sprintf("hugetlbfs-%s-%d", task.GetId(), index), memPath, ).Run() @@ -2557,11 +2590,17 @@ func (task *SGuestHotplugCpuMemTask) onGetSlotIndex(index int) { objType = "memory-backend-file" opts = map[string]string{ - "size": fmt.Sprintf("%dM", task.addMemSize), + "size": fmt.Sprintf("%dM", addMemSize), "mem-path": memPath, "share": "on", "prealloc": "on", } + + if hostNodes != nil { + opts["host-nodes"] = fmt.Sprintf("%d", *hostNodes) + opts["policy"] = "bind" + } + } else { objType = "memory-backend-ram" opts = map[string]string{ @@ -2570,14 +2609,12 @@ func (task *SGuestHotplugCpuMemTask) onGetSlotIndex(index int) { } opts["id"] = id cb := func(reason string) { - if reason == "" { - memObj := desc.NewMemDesc(objType, id, nil, nil) - memObj.Options = opts - task.memSlot = new(desc.SMemSlot) - task.memSlot.MemObj = memObj - task.memSlot.SizeMB = int64(task.addMemSize) - } - task.onAddMemObject(reason) + memObj := desc.NewMemDesc(objType, id, nil, nil) + memObj.Options = opts + memSlot := new(desc.SMemSlot) + memSlot.MemObj = memObj + memSlot.SizeMB = int64(addMemSize) + task.onAddMemObject(reason, memSlot) } task.Monitor.ObjectAdd(objType, opts, cb) } @@ -2589,7 +2626,7 @@ func (task *SGuestHotplugCpuMemTask) onAddMemFailed(reason string) { task.onFail(reason) } -func (task *SGuestHotplugCpuMemTask) onAddMemObject(reason string) { +func (task *SGuestHotplugCpuMemTask) onAddMemObject(reason string, memSlot *desc.SMemSlot) { if len(reason) > 0 { task.onAddMemFailed(reason) return @@ -2600,39 +2637,59 @@ func (task *SGuestHotplugCpuMemTask) onAddMemObject(reason string) { } cb := func(reason string) { if reason == "" { - task.memSlot.MemDev = &desc.SMemDevice{ + memSlot.MemDev = &desc.SMemDevice{ Type: "pc-dimm", Id: fmt.Sprintf("dimm%d", *task.memSlotNewIndex), } } - task.onAddMemDevice(reason) + task.onAddMemDevice(reason, memSlot) } task.Monitor.DeviceAdd("pc-dimm", params, cb) } -func (task *SGuestHotplugCpuMemTask) onAddMemDevice(reason string) { +func (task *SGuestHotplugCpuMemTask) onAddMemDevice(reason string, memSlot *desc.SMemSlot) { if len(reason) > 0 { task.onAddMemFailed(reason) return } - task.addedMemSize = task.addMemSize - task.onSucc() + task.addedMemSize = int(memSlot.SizeMB) + task.addMemSize -= int(memSlot.SizeMB) + if task.memSlots == nil { + task.memSlots = []*desc.SMemSlot{memSlot} + task.memSlotNewIndexs = []int{*task.memSlotNewIndex} + } else { + task.memSlots = append(task.memSlots, memSlot) + task.memSlotNewIndexs = append(task.memSlotNewIndexs, *task.memSlotNewIndex) + } + + if task.addMemSize > 0 { + task.startAddMem() + } else { + task.onSucc() + } } func (task *SGuestHotplugCpuMemTask) updateGuestDesc() { task.Desc.Cpu += int64(task.addedCpuCount) task.Desc.CpuDesc.Cpus += uint(task.addedCpuCount) task.Desc.Mem += int64(task.addedMemSize) + + if len(task.cpuNumaPin) > 0 { + task.Desc.CpuNumaPin = append(task.Desc.CpuNumaPin, task.cpuNumaPin...) + } + if task.addedMemSize > 0 { if task.Desc.MemDesc.MemSlots == nil { task.Desc.MemDesc.MemSlots = make([]*desc.SMemSlot, 0) } - task.Desc.MemDesc.MemSlots = append(task.Desc.MemDesc.MemSlots, task.memSlot) + task.Desc.MemDesc.MemSlots = append(task.Desc.MemDesc.MemSlots, task.memSlots...) if task.manager.numaAllocate { - hugepageId := fmt.Sprintf("%s-%d", task.getOriginId(), *task.memSlotNewIndex) - task.validateNumaAllocated(hugepageId, false, true, nil) + for i := range task.memSlotNewIndexs { + hugepageId := fmt.Sprintf("%s-%d", task.getOriginId(), task.memSlotNewIndexs[i]) + task.validateNumaAllocated(hugepageId, false, true, nil) + } } } @@ -2671,7 +2728,12 @@ func (task *SGuestHotplugCpuMemTask) onFail(reason string) { func (task *SGuestHotplugCpuMemTask) onSucc() { task.updateGuestDesc() - hostutils.TaskComplete(task.ctx, nil) + + res := jsonutils.NewDict() + if len(task.cpuNumaPin) > 0 { + res.Set("cpu_numa_pin", jsonutils.Marshal(task.Desc.CpuNumaPin)) + } + hostutils.TaskComplete(task.ctx, res) } type SGuestBlockIoThrottleTask struct { diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index 104926601e..83351fc5f2 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -153,26 +153,16 @@ func (s *SKVMGuestInstance) updateGuestDesc() error { } func (s *SKVMGuestInstance) releaseCpuNumaPin(cpuNumaPin []*desc.SCpuNumaPin) { + if !s.manager.numaAllocate { + return + } + for _, numaCpus := range cpuNumaPin { - pcpuSet, err := cpuset.Parse(*numaCpus.Pcpus) - if err != nil { - log.Errorf("failed parse %s pcpus: %s", s.GetName(), *numaCpus.Pcpus) - continue + pcpus := make([]int, 0) + for i := range numaCpus.VcpuPin { + pcpus = append(pcpus, numaCpus.VcpuPin[i].Pcpu) } - vcpuCount := int(s.Desc.Cpu) - if numaCpus.Vcpus != nil { - vcpuSet, err := cpuset.Parse(*numaCpus.Vcpus) - if err != nil { - log.Errorf("failed parse %s vcpus: %s", s.GetName(), *numaCpus.Vcpus) - continue - } - vcpuCount = vcpuSet.Size() - } - hostNodes := -1 - if numaCpus.HostNodes != nil { - hostNodes = int(*numaCpus.HostNodes) - } - s.manager.cpuSet.ReleaseNumaCpus(numaCpus.SizeMB, hostNodes, pcpuSet.ToSlice(), vcpuCount) + s.manager.cpuSet.ReleaseNumaCpus(numaCpus.SizeMB, int(*numaCpus.NodeId), pcpus, len(numaCpus.VcpuPin)) } } @@ -209,12 +199,15 @@ func (s *SKVMGuestInstance) reallocateMigrateNumaNodes() error { if len(nodeNumaCpus) > 0 { for nodeId, numaCpus := range nodeNumaCpus { unodeId := uint16(nodeId) - pcpus := cpuset.NewCPUSet(numaCpus.Cpuset...).String() + vcpuPin := make([]desc.SVCpuPin, len(numaCpus.Cpuset)) + for i := range numaCpus.Cpuset { + vcpuPin[i].Pcpu = numaCpus.Cpuset[i] + } memPin := &desc.SCpuNumaPin{ SizeMB: numaCpus.MemSizeKB / 1024, // MB - HostNodes: &unodeId, - Pcpus: &pcpus, - Regular: numaCpus.Regular, + NodeId: &unodeId, + VcpuPin: vcpuPin, + Unregular: numaCpus.Unregular, } cpuNumaPin = append(cpuNumaPin, memPin) } @@ -222,9 +215,9 @@ func (s *SKVMGuestInstance) reallocateMigrateNumaNodes() error { if len(cpuNumaPin) > 0 { s.Desc.CpuNumaPin = cpuNumaPin - s.Desc.MemDesc.Mem.SMemDesc.SetHostNodes(int(*cpuNumaPin[0].HostNodes)) + s.Desc.MemDesc.Mem.SMemDesc.SetHostNodes(int(*cpuNumaPin[0].NodeId)) for i := range s.Desc.MemDesc.Mem.Mems { - s.Desc.MemDesc.Mem.Mems[i].SetHostNodes(int(*cpuNumaPin[i+1].HostNodes)) + s.Desc.MemDesc.Mem.Mems[i].SetHostNodes(int(*cpuNumaPin[i+1].NodeId)) } } else { s.Desc.MemDesc.Mem.SMemDesc.SetHostNodes(-1) @@ -236,11 +229,13 @@ func (s *SKVMGuestInstance) reallocateMigrateNumaNodes() error { return nil } -func (s *SKVMGuestInstance) validateNumaAllocated(keywords string, isMigrate, isHotPlug bool, vcpuOrder []string) error { +func (s *SKVMGuestInstance) validateNumaAllocated(keywords string, isMigrate, isHotPlug bool, vcpuOrder [][]int) error { if len(s.Desc.CpuNumaPin) > 0 { if isMigrate { for i := range s.Desc.CpuNumaPin { - s.Desc.CpuNumaPin[i].Vcpus = &vcpuOrder[i] + for j := range s.Desc.CpuNumaPin[i].VcpuPin { + s.Desc.CpuNumaPin[i].VcpuPin[j].Vcpu = vcpuOrder[i][j] + } } return SaveLiveDesc(s, s.Desc) } @@ -308,12 +303,16 @@ func (s *SKVMGuestInstance) validateNumaAllocated(keywords string, isMigrate, is var cpuNumaPin = make([]*desc.SCpuNumaPin, 0) for nodeId, numaCpus := range nodeNumaCpus { unodeId := uint16(nodeId) - pcpus := cpuset.NewCPUSet(numaCpus.Cpuset...).String() + vcpuPin := make([]desc.SVCpuPin, len(numaCpus.Cpuset)) + for i := range numaCpus.Cpuset { + vcpuPin[i].Pcpu = numaCpus.Cpuset[i] + } + memPin := &desc.SCpuNumaPin{ SizeMB: numaCpus.MemSizeKB / 1024, // MB - HostNodes: &unodeId, - Pcpus: &pcpus, - Regular: numaCpus.Regular, + NodeId: &unodeId, + VcpuPin: vcpuPin, + Unregular: numaCpus.Unregular, } cpuNumaPin = append(cpuNumaPin, memPin) } @@ -325,7 +324,9 @@ func (s *SKVMGuestInstance) validateNumaAllocated(keywords string, isMigrate, is if len(vcpuOrder) > 0 { for i := range cpuNumaPin { - cpuNumaPin[i].Vcpus = &vcpuOrder[i] + for j := range cpuNumaPin[i].VcpuPin { + cpuNumaPin[i].VcpuPin[j].Vcpu = vcpuOrder[i][j] + } } } @@ -378,66 +379,79 @@ func (s *SKVMGuestInstance) initLiveDescFromSourceGuest(srcDesc *desc.SGuestDesc srcDesc.Nics[i].DownscriptPath = s.getNicDownScriptPath(srcDesc.Nics[i]) } - nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(srcDesc.Cpu), srcDesc.Mem*1024, len(srcDesc.MemDesc.Mem.Mems)+1) - if err != nil { - return errors.Wrap(err, "AllocCpusetWithNodeCount") - } - - var cpus = make([]int, 0) - var cpuNumaPin = make([]*desc.SCpuNumaPin, 0) - for nodeId, numaCpus := range nodeNumaCpus { - if s.manager.numaAllocate { - unodeId := uint16(nodeId) - pcpus := cpuset.NewCPUSet(numaCpus.Cpuset...).String() - memPin := &desc.SCpuNumaPin{ - SizeMB: numaCpus.MemSizeKB / 1024, // MB - HostNodes: &unodeId, - Pcpus: &pcpus, - Regular: numaCpus.Regular, - } - cpuNumaPin = append(cpuNumaPin, memPin) - } - cpus = append(cpus, numaCpus.Cpuset...) - } - - if s.manager.numaAllocate { - srcDesc.VcpuPin = nil - srcDesc.CpuNumaPin = nil + var cpuNumaPin []*desc.SCpuNumaPin + if len(s.Desc.CpuNumaPin) > 0 { + // cpu numa pin allocated by controller + cpuNumaPin = s.Desc.CpuNumaPin } else { - if scpuset, ok := srcDesc.Metadata[api.VM_METADATA_CGROUP_CPUSET]; ok { - s.manager.cpuSet.Lock.Lock() - s.manager.cpuSet.ReleaseCpus(cpus, int(srcDesc.Cpu)) - s.manager.cpuSet.Lock.Unlock() - - cpusetJson, err := jsonutils.ParseString(scpuset) - if err != nil { - log.Errorf("failed parse server %s cpuset %s: %s", s.Id, scpuset, err) - return errors.Errorf("failed parse server %s cpuset %s: %s", s.Id, scpuset, err) - } - input := new(api.ServerCPUSetInput) - err = cpusetJson.Unmarshal(input) - if err != nil { - log.Errorf("failed unmarshal server %s cpuset %s", s.Id, err) - return errors.Errorf("failed unmarshal server %s cpuset %s", s.Id, err) - } - cpus = input.CPUS + // allocate cpu numa pin local + nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(srcDesc.Cpu), srcDesc.Mem*1024, len(srcDesc.MemDesc.Mem.Mems)+1) + if err != nil { + return errors.Wrap(err, "AllocCpusetWithNodeCount") } - srcDesc.VcpuPin = []desc.SCpuPin{ - { - Vcpus: fmt.Sprintf("0-%d", srcDesc.Cpu-1), - Pcpus: cpuset.NewCPUSet(cpus...).String(), - }, + var cpus = make([]int, 0) + cpuNumaPin = make([]*desc.SCpuNumaPin, 0) + for nodeId, numaCpus := range nodeNumaCpus { + if s.manager.numaAllocate { + unodeId := uint16(nodeId) + vcpuPin := make([]desc.SVCpuPin, len(numaCpus.Cpuset)) + for i := range numaCpus.Cpuset { + vcpuPin[i].Pcpu = numaCpus.Cpuset[i] + } + + memPin := &desc.SCpuNumaPin{ + SizeMB: numaCpus.MemSizeKB / 1024, // MB + NodeId: &unodeId, + VcpuPin: vcpuPin, + Unregular: numaCpus.Unregular, + } + cpuNumaPin = append(cpuNumaPin, memPin) + } + cpus = append(cpus, numaCpus.Cpuset...) } - for i := range srcDesc.CpuNumaPin { - srcDesc.CpuNumaPin[i].Regular = false + + if s.manager.numaAllocate { + // reset origin cpu numa pin + srcDesc.VcpuPin = nil + srcDesc.CpuNumaPin = nil + } else { + // if host not enable cpu numa pin + if scpuset, ok := srcDesc.Metadata[api.VM_METADATA_CGROUP_CPUSET]; ok { + s.manager.cpuSet.Lock.Lock() + s.manager.cpuSet.ReleaseCpus(cpus, int(srcDesc.Cpu)) + s.manager.cpuSet.Lock.Unlock() + + cpusetJson, err := jsonutils.ParseString(scpuset) + if err != nil { + log.Errorf("failed parse server %s cpuset %s: %s", s.Id, scpuset, err) + return errors.Errorf("failed parse server %s cpuset %s: %s", s.Id, scpuset, err) + } + input := new(api.ServerCPUSetInput) + err = cpusetJson.Unmarshal(input) + if err != nil { + log.Errorf("failed unmarshal server %s cpuset %s", s.Id, err) + return errors.Errorf("failed unmarshal server %s cpuset %s", s.Id, err) + } + cpus = input.CPUS + } + + srcDesc.VcpuPin = []desc.SCpuPin{ + { + Vcpus: fmt.Sprintf("0-%d", srcDesc.Cpu-1), + Pcpus: cpuset.NewCPUSet(cpus...).String(), + }, + } + for i := range srcDesc.CpuNumaPin { + srcDesc.CpuNumaPin[i].Unregular = true + } } } if len(cpuNumaPin) > 0 { - srcDesc.MemDesc.Mem.SMemDesc.SetHostNodes(int(*cpuNumaPin[0].HostNodes)) + srcDesc.MemDesc.Mem.SMemDesc.SetHostNodes(int(*cpuNumaPin[0].NodeId)) for i := range srcDesc.MemDesc.Mem.Mems { - srcDesc.MemDesc.Mem.Mems[i].SetHostNodes(int(*cpuNumaPin[i+1].HostNodes)) + srcDesc.MemDesc.Mem.Mems[i].SetHostNodes(int(*cpuNumaPin[i+1].NodeId)) } srcDesc.CpuNumaPin = cpuNumaPin } else { @@ -448,7 +462,7 @@ func (s *SKVMGuestInstance) initLiveDescFromSourceGuest(srcDesc *desc.SGuestDesc } s.Desc = srcDesc - err = s.loadGuestPciAddresses() + err := s.loadGuestPciAddresses() if err != nil { return errors.Wrap(err, "initLiveDescFromSourceGuest") } @@ -658,26 +672,12 @@ func (s *SKVMGuestInstance) loadGuestCpuset(m *SGuestManager) error { m.cpuSet.LoadCpus(pcpuSet.ToSlice(), vcpuSet.Size()) } for _, numaCpuset := range s.Desc.CpuNumaPin { - pcpuSet, err := cpuset.Parse(*numaCpuset.Pcpus) - if err != nil { - log.Errorf("failed parse %s pcpus: %s", s.GetName(), *numaCpuset.Pcpus) - continue - } - vcpuCount := int(s.Desc.Cpu) - if numaCpuset.Vcpus != nil { - vcpuSet, err := cpuset.Parse(*numaCpuset.Vcpus) - if err != nil { - log.Errorf("failed parse %s vcpus: %s", s.GetName(), *numaCpuset.Vcpus) - continue - } - vcpuCount = vcpuSet.Size() - } - hostNodes := -1 - if numaCpuset.HostNodes != nil { - hostNodes = int(*numaCpuset.HostNodes) + pcpus := make([]int, 0) + for i := range numaCpuset.VcpuPin { + pcpus = append(pcpus, numaCpuset.VcpuPin[i].Pcpu) } - m.cpuSet.LoadNumaCpus(numaCpuset.SizeMB, hostNodes, pcpuSet.ToSlice(), vcpuCount) + m.cpuSet.LoadNumaCpus(numaCpuset.SizeMB, int(*numaCpuset.NodeId), pcpus, len(numaCpuset.VcpuPin)) } } return nil @@ -748,7 +748,7 @@ func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interfa return nil, errors.Wrap(err, "fuse mount") } - var vcpuOrder = make([]string, 0) + var vcpuOrder = make([][]int, 0) isMigrate := jsonutils.QueryBoolean(data, "need_migrate", false) if isMigrate { var sourceDesc = new(desc.SGuestDesc) @@ -757,10 +757,11 @@ func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interfa return nil, errors.Wrap(err, "unmarshal src desc") } for i := range sourceDesc.CpuNumaPin { - if sourceDesc.CpuNumaPin[i].Vcpus != nil { - vcpus := *sourceDesc.CpuNumaPin[i].Vcpus - vcpuOrder = append(vcpuOrder, vcpus) + vcpus := make([]int, 0) + for j := range sourceDesc.CpuNumaPin[i].VcpuPin { + vcpus = append(vcpus, sourceDesc.CpuNumaPin[i].VcpuPin[j].Vcpu) } + vcpuOrder = append(vcpuOrder, vcpus) } err = s.initLiveDescFromSourceGuest(sourceDesc) } else { @@ -2719,9 +2720,13 @@ func (s *SKVMGuestInstance) setCgroupCPUSet() error { } cpusetStr = strings.Join(cpus, ",") } else { + pcpuSetBuilder := cpuset.NewBuilder() for i := range s.Desc.CpuNumaPin { - cpusetStr += fmt.Sprintf(",%s", *s.Desc.CpuNumaPin[i].Pcpus) + for j := range s.Desc.CpuNumaPin[i].VcpuPin { + pcpuSetBuilder.Add(s.Desc.CpuNumaPin[i].VcpuPin[j].Pcpu) + } } + cpusetStr = pcpuSetBuilder.Result().String() } guestPid := strconv.Itoa(s.GetPid()) @@ -2737,20 +2742,18 @@ func (s *SKVMGuestInstance) setCgroupCPUSet() error { } for i := range s.Desc.CpuNumaPin { - if !s.Desc.CpuNumaPin[i].Regular || s.Desc.CpuNumaPin[i].Vcpus == nil { + if s.Desc.CpuNumaPin[i].Unregular || s.Desc.CpuNumaPin[i].VcpuPin == nil { continue } - vcpuSet, _ := cpuset.Parse(*s.Desc.CpuNumaPin[i].Vcpus) - pcpuSet := *s.Desc.CpuNumaPin[i].Pcpus - for _, vcpuId := range vcpuSet.ToSlice() { - vcpuThreadId, ok := vcpuThreads[vcpuId] + for j := range s.Desc.CpuNumaPin[i].VcpuPin { + vcpuThreadId, ok := vcpuThreads[s.Desc.CpuNumaPin[i].VcpuPin[j].Vcpu] if !ok { - return errors.Errorf("failed get vcpu %d thread id from %v", vcpuId, vcpuThreads) + return errors.Errorf("failed get vcpu %d thread id from %v", s.Desc.CpuNumaPin[i].VcpuPin[j].Vcpu, vcpuThreads) } - + pcpu := s.Desc.CpuNumaPin[i].VcpuPin[j].Pcpu vcpuCgname := path.Join(cgName, vcpuThreadId) - taskVcpu := cgrouputils.NewCGroupSubCPUSetTask(guestPid, vcpuCgname, 0, pcpuSet, []string{vcpuThreadId}) + taskVcpu := cgrouputils.NewCGroupSubCPUSetTask(guestPid, vcpuCgname, 0, strconv.Itoa(pcpu), []string{vcpuThreadId}) if !taskVcpu.SetTask() { return errors.Errorf("Vcpu set cgroup cpuset task failed") } @@ -2778,12 +2781,16 @@ func (s *SKVMGuestInstance) allocGuestNumaCpuset() error { for nodeId, numaCpus := range nodeNumaCpus { if s.manager.numaAllocate { unodeId := uint16(nodeId) - pcpus := cpuset.NewCPUSet(numaCpus.Cpuset...).String() + vcpuPin := make([]desc.SVCpuPin, len(numaCpus.Cpuset)) + for i := range numaCpus.Cpuset { + vcpuPin[i].Pcpu = numaCpus.Cpuset[i] + } + memPin := &desc.SCpuNumaPin{ SizeMB: numaCpus.MemSizeKB / 1024, // MB - HostNodes: &unodeId, - Pcpus: &pcpus, - Regular: numaCpus.Regular, + NodeId: &unodeId, + VcpuPin: vcpuPin, + Unregular: numaCpus.Unregular, } cpuNumaPin = append(cpuNumaPin, memPin) } @@ -2996,11 +3003,9 @@ func (s *SKVMGuestInstance) hotPlugCpus() error { var vcpuSet = make([]int, 0) if len(s.Desc.MemDesc.Mem.Mems) > 0 { for i := range s.Desc.CpuNumaPin { - vcpus, err := cpuset.Parse(*s.Desc.CpuNumaPin[i].Vcpus) - if err != nil { - return errors.Wrap(err, "parse vcpus") + for j := range s.Desc.CpuNumaPin[i].VcpuPin { + vcpuSet = append(vcpuSet, s.Desc.CpuNumaPin[i].VcpuPin[j].Vcpu) } - vcpuSet = append(vcpuSet, vcpus.ToSlice()...) } } return s.startHotPlugVcpus(vcpuSet) diff --git a/pkg/hostman/guestman/qemu-kvmhelper.go b/pkg/hostman/guestman/qemu-kvmhelper.go index 2e4e2fa2d4..b027622e23 100644 --- a/pkg/hostman/guestman/qemu-kvmhelper.go +++ b/pkg/hostman/guestman/qemu-kvmhelper.go @@ -883,11 +883,13 @@ func (s *SKVMGuestInstance) initCpuDesc(cpuMax uint) error { } s.Desc.CpuDesc = cpuDesc - err = s.allocGuestNumaCpuset() - if err != nil { - return err + // if region not allocate cpu numa pin + if len(s.Desc.CpuNumaPin) == 0 { + err = s.allocGuestNumaCpuset() + if err != nil { + return err + } } - return nil } @@ -917,9 +919,15 @@ func (s *SKVMGuestInstance) initGuestMemObjects(memSizeMB int64) error { var numaMems int64 var numaCpus = int(s.Desc.CpuDesc.MaxCpus) / len(s.Desc.CpuNumaPin) var leastCpus = int(s.Desc.CpuDesc.MaxCpus) % len(s.Desc.CpuNumaPin) - var numaVcpuCount = int(s.Desc.Cpu) / len(s.Desc.CpuNumaPin) + var cpuStart = 0 + var cpuEnd = numaCpus - 1 + var mems = make([]desc.SMemDesc, 0) for i := 0; i < len(s.Desc.CpuNumaPin); i++ { + if s.Desc.CpuNumaPin[i].SizeMB <= 0 { + continue + } + numaMems += s.Desc.CpuNumaPin[i].SizeMB memId := "mem" nodeId := uint16(i) @@ -927,25 +935,23 @@ func (s *SKVMGuestInstance) initGuestMemObjects(memSizeMB int64) error { memId += strconv.Itoa(i - 1) } - vcpuCount := numaVcpuCount if i == 0 { - vcpuCount += int(s.Desc.Cpu) % len(s.Desc.CpuNumaPin) - } - - cpuStart := i * numaCpus - cpuEnd := (i+1)*numaCpus - 1 - if i+1 == len(s.Desc.CpuNumaPin) { cpuEnd += leastCpus } vcpus := fmt.Sprintf("%d-%d", cpuStart, cpuEnd) - vcpuAlloc := fmt.Sprintf("%d-%d", cpuStart, cpuStart+vcpuCount-1) - s.Desc.CpuNumaPin[i].Vcpus = &vcpuAlloc - if !s.Desc.CpuNumaPin[i].Regular { + for j := range s.Desc.CpuNumaPin[i].VcpuPin { + s.Desc.CpuNumaPin[i].VcpuPin[j].Vcpu = cpuStart + j + } + + cpuStart = cpuEnd + 1 + cpuEnd = cpuStart + numaCpus - 1 + + if s.Desc.CpuNumaPin[i].Unregular { continue } memDesc := desc.NewMemDesc(s.memObjectType(), memId, &nodeId, &vcpus) - memDesc.Options = s.getMemObjectOptions(s.Desc.CpuNumaPin[i].SizeMB, s.Desc.Uuid, s.Desc.CpuNumaPin[i].HostNodes) + memDesc.Options = s.getMemObjectOptions(s.Desc.CpuNumaPin[i].SizeMB, s.Desc.Uuid, s.Desc.CpuNumaPin[i].NodeId) mems = append(mems, *memDesc) } if len(mems) == 0 { @@ -954,9 +960,6 @@ func (s *SKVMGuestInstance) initGuestMemObjects(memSizeMB int64) error { return nil } - if numaMems != memSizeMB { - return errors.Errorf("numa memory size not equal request mem size") - } s.Desc.MemDesc.Mem = desc.NewMemsDesc(mems[0], mems[1:]) return nil } diff --git a/pkg/hostman/hostinfo/hostinfo.go b/pkg/hostman/hostinfo/hostinfo.go index 80b737f6a0..9f6d664b10 100644 --- a/pkg/hostman/hostinfo/hostinfo.go +++ b/pkg/hostman/hostinfo/hostinfo.go @@ -97,6 +97,7 @@ type SHostInfo struct { onHostDown string reservedCpusInfo *api.HostReserveCpusInput enableNumaAllocate bool + cpuCmtBound int IsolatedDeviceMan isolated_device.IsolatedDeviceManager @@ -541,6 +542,9 @@ func (h *SHostInfo) detectHostInfo() error { h.sysinfo.Topology = topoInfo h.sysinfo.CPUInfo = cpuInfo + if err = h.GetNodeHugepages(); err != nil { + return errors.Wrap(err, "GetNodeHugepages") + } system_service.Init() if options.HostOptions.CheckSystemServices { if err := h.checkSystemServices(); err != nil { @@ -630,6 +634,33 @@ func (h *SHostInfo) EnableTransparentHugepages() { } } +func (h *SHostInfo) GetNodeHugepages() error { + if options.HostOptions.HugepagesOption != "native" { + return nil + } + + hugepageSizeKB := h.sysinfo.HugepageSizeKb + nodeHugepages := make([]hostapi.HostNodeHugepageNr, len(h.sysinfo.Topology.Nodes)) + + for i := range h.sysinfo.Topology.Nodes { + nodeId := h.sysinfo.Topology.Nodes[i].ID + nodeHugepagePath := fmt.Sprintf("/sys/devices/system/node/node%d/hugepages/hugepages-%dkB", nodeId, hugepageSizeKB) + if !fileutils2.Exists(nodeHugepagePath) { + return errors.Errorf("node %s has no hugepages ?", nodeHugepagePath) + } + nrHugepage, err := fileutils2.FileGetIntContent(path.Join(nodeHugepagePath, "nr_hugepages")) + if err != nil { + return errors.Wrap(err, "get node nr hugepage") + } + + nodeHugepages[i].NodeId = nodeId + nodeHugepages[i].HugepageNr = nrHugepage + } + + h.sysinfo.NodeHugepages = nodeHugepages + return nil +} + func (h *SHostInfo) GetMemory() int { return h.Mem.Total } @@ -1167,6 +1198,7 @@ func (h *SHostInfo) initHostRecord() (*api.HostDetails, error) { } h.HostId = hostInfo.Id + h.cpuCmtBound = int(hostInfo.CpuCmtbound) hostInfo, err = h.updateHostMetadata(hostInfo.Name) if err != nil { return nil, errors.Wrap(err, "updateHostMetadata") @@ -2466,6 +2498,10 @@ func (h *SHostInfo) GetContainerRuntimeEndpoint() string { return options.HostOptions.ContainerRuntimeEndpoint } +func (h *SHostInfo) CpuCmtBound() int { + return h.cpuCmtBound +} + func NewHostInfo() (*SHostInfo, error) { var res = new(SHostInfo) res.sysinfo = &SSysInfo{} diff --git a/pkg/hostman/hostinfo/hostinfohelper.go b/pkg/hostman/hostinfo/hostinfohelper.go index 05258f700e..ee0dd63ab2 100644 --- a/pkg/hostman/hostinfo/hostinfohelper.go +++ b/pkg/hostman/hostinfo/hostinfohelper.go @@ -327,9 +327,10 @@ type SSysInfo struct { StorageType string `json:"storage_type"` - HugepagesOption string `json:"hugepages_option"` - HugepageSizeKb int `json:"hugepage_size_kb"` - HugepageNr *int `json:"hugepage_nr"` + HugepagesOption string `json:"hugepages_option"` + HugepageSizeKb int `json:"hugepage_size_kb"` + HugepageNr *int `json:"hugepage_nr"` + NodeHugepages []hostapi.HostNodeHugepageNr `json:"node_hugepages"` Topology *hostapi.HostTopology `json:"topology"` CPUInfo *hostapi.HostCPUInfo `json:"cpu_info"` diff --git a/pkg/hostman/hostutils/hostutils.go b/pkg/hostman/hostutils/hostutils.go index 547f249170..1fc356ad82 100644 --- a/pkg/hostman/hostutils/hostutils.go +++ b/pkg/hostman/hostutils/hostutils.go @@ -60,6 +60,7 @@ type IHost interface { IsHugepagesEnabled() bool HugepageSizeKb() int IsNumaAllocateEnabled() bool + CpuCmtBound() int IsKvmSupport() bool IsNestedVirtualization() bool diff --git a/pkg/hostman/options/options.go b/pkg/hostman/options/options.go index f32fcfc434..7ce45e2df5 100644 --- a/pkg/hostman/options/options.go +++ b/pkg/hostman/options/options.go @@ -132,8 +132,9 @@ type SHostOptions struct { SetVncPassword bool `default:"true" help:"Auto set vnc password after monitor connected"` UseBootVga bool `default:"false" help:"Use boot VGA GPU for guest"` - EnableCpuBinding bool `default:"true" help:"Enable cpu binding and rebalance"` - EnableOpenflowController bool `default:"false"` + EnableHostAgentNumaAllocate bool `default:"false" help:"Enable host agent numa allocate"` + EnableCpuBinding bool `default:"true" help:"Enable cpu binding and rebalance"` + EnableOpenflowController bool `default:"false"` PingRegionInterval int `default:"60" help:"interval to ping region, deefault is 1 minute"` LogSystemdUnits []string `help:"Systemd units log collected by fluent-bit"` diff --git a/pkg/scheduler/algorithm/predicates/class_metadata_predicate.go b/pkg/scheduler/algorithm/predicates/class_metadata_predicate.go index 73db6c3beb..564f2e73bb 100644 --- a/pkg/scheduler/algorithm/predicates/class_metadata_predicate.go +++ b/pkg/scheduler/algorithm/predicates/class_metadata_predicate.go @@ -62,6 +62,10 @@ func (p *ClassMetadataPredicate) Clone() core.FitPredicate { func (p *ClassMetadataPredicate) PreExecute(ctx context.Context, u *core.Unit, cs []core.Candidater) (bool, error) { info := u.SchedData() + if info.ResetCpuNumaPin { + return false, nil + } + tenant, err := db.TenantCacheManager.FetchTenantById(ctx, info.Project) if err != nil { return false, errors.Wrapf(err, "unable to fetch tenant %s", info.Project) diff --git a/pkg/scheduler/algorithm/predicates/cloudprovider_schedtag_predicate.go b/pkg/scheduler/algorithm/predicates/cloudprovider_schedtag_predicate.go index 13a073a721..acec420a6b 100644 --- a/pkg/scheduler/algorithm/predicates/cloudprovider_schedtag_predicate.go +++ b/pkg/scheduler/algorithm/predicates/cloudprovider_schedtag_predicate.go @@ -48,6 +48,9 @@ func (p *CloudproviderSchedtagPredicate) PreExecute(ctx context.Context, u *core if driver == nil || !driver.DoScheduleCloudproviderTagFilter() { return false, nil } + if u.SchedData().ResetCpuNumaPin { + return false, nil + } return p.ServerBaseSchedtagPredicate.PreExecute(ctx, u, cs) } diff --git a/pkg/scheduler/algorithm/predicates/guest/image_predicate.go b/pkg/scheduler/algorithm/predicates/guest/image_predicate.go index c47bda2cab..65483a625f 100644 --- a/pkg/scheduler/algorithm/predicates/guest/image_predicate.go +++ b/pkg/scheduler/algorithm/predicates/guest/image_predicate.go @@ -43,6 +43,10 @@ func (f *ImagePredicate) Clone() core.FitPredicate { } func (f *ImagePredicate) PreExecute(ctx context.Context, u *core.Unit, cs []core.Candidater) (bool, error) { + if u.SchedData().ResetCpuNumaPin { + return false, nil + } + disks := u.SchedData().Disks if len(disks) == 0 { return false, nil diff --git a/pkg/scheduler/algorithm/predicates/guest/memory_predicate.go b/pkg/scheduler/algorithm/predicates/guest/memory_predicate.go index ffb6f1a64e..4ad0ae5ec2 100644 --- a/pkg/scheduler/algorithm/predicates/guest/memory_predicate.go +++ b/pkg/scheduler/algorithm/predicates/guest/memory_predicate.go @@ -16,7 +16,11 @@ package guest import ( "context" + "fmt" + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/apis/scheduler" "yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates" "yunion.io/x/onecloud/pkg/scheduler/core" ) @@ -64,6 +68,40 @@ func (p *MemoryPredicate) Execute(ctx context.Context, u *core.Unit, c core.Cand h.AppendInsufficientResourceError(reqMemSize, totalMemSize, freeMemSize) } + if cpuNumaFree := getter.GetFreeCpuNuma(); cpuNumaFree != nil { + allcateEnough := false + reqCpuCount := d.Ncpu + if d.CpuNumaPin != nil { + nodeCount := len(d.CpuNumaPin) + if scheduler.NodesFreeCpuEnough(nodeCount, d.Ncpu, cpuNumaFree) && + scheduler.NodesFreeMemSizeEnough(nodeCount, int(reqMemSize), cpuNumaFree) { + allcateEnough = true + } + } else { + for nodeCount := 1; nodeCount <= len(cpuNumaFree); nodeCount *= 2 { + if nodeCount > reqCpuCount { + break + } + + if !scheduler.NodesFreeCpuEnough(nodeCount, d.Ncpu, cpuNumaFree) { + continue + } + + if !scheduler.NodesFreeMemSizeEnough(nodeCount, int(reqMemSize), cpuNumaFree) { + continue + } + allcateEnough = true + } + } + + if !allcateEnough { + h.AppendPredicateFailMsg( + fmt.Sprintf("cpu numa free %s can't alloc with req mem %v req cpu %v ", + jsonutils.Marshal(cpuNumaFree).String(), d.Memory, d.Ncpu), + ) + } + } + h.SetCapacity(freeMemSize / reqMemSize) return h.GetResult() } diff --git a/pkg/scheduler/algorithm/predicates/guest/migrate_predicate.go b/pkg/scheduler/algorithm/predicates/guest/migrate_predicate.go index 66274f080b..893e12ef9a 100644 --- a/pkg/scheduler/algorithm/predicates/guest/migrate_predicate.go +++ b/pkg/scheduler/algorithm/predicates/guest/migrate_predicate.go @@ -36,6 +36,10 @@ func (p *MigratePredicate) Clone() core.FitPredicate { } func (p *MigratePredicate) PreExecute(ctx context.Context, u *core.Unit, cs []core.Candidater) (bool, error) { + if u.SchedData().ResetCpuNumaPin { + return false, nil + } + return len(u.SchedData().HostId) > 0, nil } diff --git a/pkg/scheduler/algorithm/predicates/guest/storage_predicate.go b/pkg/scheduler/algorithm/predicates/guest/storage_predicate.go index 1f095617ca..617fb004ad 100644 --- a/pkg/scheduler/algorithm/predicates/guest/storage_predicate.go +++ b/pkg/scheduler/algorithm/predicates/guest/storage_predicate.go @@ -50,6 +50,10 @@ func (p *StoragePredicate) PreExecute(ctx context.Context, u *core.Unit, cs []co if driver != nil && !driver.DoScheduleStorageFilter() { return false, nil } + if u.SchedData().ResetCpuNumaPin { + return false, nil + } + return true, nil } diff --git a/pkg/scheduler/algorithm/predicates/isolated_device_predicate.go b/pkg/scheduler/algorithm/predicates/isolated_device_predicate.go index f83154f813..6fbae94116 100644 --- a/pkg/scheduler/algorithm/predicates/isolated_device_predicate.go +++ b/pkg/scheduler/algorithm/predicates/isolated_device_predicate.go @@ -38,6 +38,11 @@ func (f *IsolatedDevicePredicate) Clone() core.FitPredicate { func (f *IsolatedDevicePredicate) PreExecute(ctx context.Context, u *core.Unit, cs []core.Candidater) (bool, error) { data := u.SchedData() + + if data.ResetCpuNumaPin { + return false, nil + } + if len(data.IsolatedDevices) > 0 { return true, nil } diff --git a/pkg/scheduler/algorithm/predicates/network_predicate.go b/pkg/scheduler/algorithm/predicates/network_predicate.go index c6f390d678..3b34c9c08f 100644 --- a/pkg/scheduler/algorithm/predicates/network_predicate.go +++ b/pkg/scheduler/algorithm/predicates/network_predicate.go @@ -75,6 +75,11 @@ type INetworkNicCountGetter interface { func (p *NetworkPredicate) PreExecute(ctx context.Context, u *core.Unit, cs []core.Candidater) (bool, error) { data := u.SchedData() + + if data.ResetCpuNumaPin { + return false, nil + } + if len(data.Networks) == 0 { return false, nil } diff --git a/pkg/scheduler/algorithm/predicates/predicates.go b/pkg/scheduler/algorithm/predicates/predicates.go index ac510f79e5..ffa3b98b1a 100644 --- a/pkg/scheduler/algorithm/predicates/predicates.go +++ b/pkg/scheduler/algorithm/predicates/predicates.go @@ -417,6 +417,10 @@ func (p *BaseSchedtagPredicate) PreExecute(ctx context.Context, sp ISchedtagPred return false, nil } + if u.SchedData().ResetCpuNumaPin { + return false, nil + } + p.Hypervisor = u.GetHypervisor() p.Provider = u.SchedInfo.Provider diff --git a/pkg/scheduler/algorithm/predicates/sku_predicate.go b/pkg/scheduler/algorithm/predicates/sku_predicate.go index d3b2e11502..1a8fcde40b 100644 --- a/pkg/scheduler/algorithm/predicates/sku_predicate.go +++ b/pkg/scheduler/algorithm/predicates/sku_predicate.go @@ -36,6 +36,10 @@ func (p *InstanceTypePredicate) Clone() core.FitPredicate { func (p *InstanceTypePredicate) PreExecute(ctx context.Context, u *core.Unit, cs []core.Candidater) (bool, error) { driver := u.GetHypervisorDriver() + if u.SchedData().ResetCpuNumaPin { + return false, nil + } + if u.SchedData().InstanceType == "" || (driver == nil || !driver.DoScheduleSKUFilter()) { return false, nil } diff --git a/pkg/scheduler/algorithm/priorities/guest/cpunumapin.go b/pkg/scheduler/algorithm/priorities/guest/cpunumapin.go new file mode 100644 index 0000000000..3aec3a28ee --- /dev/null +++ b/pkg/scheduler/algorithm/priorities/guest/cpunumapin.go @@ -0,0 +1,62 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package guest + +import ( + "yunion.io/x/onecloud/pkg/apis/scheduler" + "yunion.io/x/onecloud/pkg/scheduler/algorithm/priorities" + "yunion.io/x/onecloud/pkg/scheduler/core" +) + +type CpuNumaPinPriority struct { + priorities.BasePriority +} + +func (p *CpuNumaPinPriority) Name() string { + return "cpu_numa_pin" +} + +func (p *CpuNumaPinPriority) Clone() core.Priority { + return &CpuNumaPinPriority{} +} + +func (p *CpuNumaPinPriority) Map(u *core.Unit, c core.Candidater) (core.HostPriority, error) { + h := priorities.NewPriorityHelper(p, u, c) + + getter := c.Getter() + if getter.Host().EnableNumaAllocate && getter.NumaAllocateEnabled() && len(u.SchedData().CpuNumaPin) == 0 { + cpuNumaFree := getter.GetFreeCpuNuma() + + reqCpuCount := u.SchedInfo.Ncpu + reqMemSize := u.SchedInfo.Memory + nodeCount := 1 + for ; nodeCount <= len(cpuNumaFree); nodeCount *= 2 { + if nodeCount > reqCpuCount { + nodeCount = 0 + break + } + + if scheduler.NodesFreeCpuEnough(nodeCount, reqCpuCount, cpuNumaFree) && + scheduler.NodesFreeMemSizeEnough(nodeCount, int(reqMemSize), cpuNumaFree) { + break + } + } + if nodeCount > 0 && nodeCount <= len(cpuNumaFree) { + score := 100 / nodeCount + h.SetScore(int(score)) + } + } + return h.GetResult() +} diff --git a/pkg/scheduler/algorithmprovider/defaults.go b/pkg/scheduler/algorithmprovider/defaults.go index dbd99945c9..f6d225a73d 100644 --- a/pkg/scheduler/algorithmprovider/defaults.go +++ b/pkg/scheduler/algorithmprovider/defaults.go @@ -59,5 +59,6 @@ func defaultPriorities() sets.String { factory.RegisterPriority("guest-lowload", &priorityguest.LowLoadPriority{}, 1), factory.RegisterPriority("guest-creating", &priorityguest.CreatingPriority{}, 1), factory.RegisterPriority("guest-capacity", &priorityguest.CapacityPriority{}, 1), + factory.RegisterPriority("guest-cpunumapin", &priorityguest.CpuNumaPinPriority{}, 1), ) } diff --git a/pkg/scheduler/cache/candidate/baremetals.go b/pkg/scheduler/cache/candidate/baremetals.go index 3b8c1b1169..f746c69efc 100644 --- a/pkg/scheduler/cache/candidate/baremetals.go +++ b/pkg/scheduler/cache/candidate/baremetals.go @@ -19,6 +19,7 @@ import ( "yunion.io/x/sqlchemy" computeapi "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/apis/scheduler" "yunion.io/x/onecloud/pkg/compute/baremetal" computemodels "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/scheduler/core" @@ -44,6 +45,14 @@ func (h baremetalGetter) FreeMemorySize(_ bool) int64 { return h.bm.FreeMemSize() } +func (h *baremetalGetter) NumaAllocateEnabled() bool { + return false +} + +func (h *baremetalGetter) GetFreeCpuNuma() []*scheduler.SFreeNumaCpuMem { + return nil +} + func (h baremetalGetter) IsEmpty() bool { return h.bm.ServerID == "" } @@ -106,6 +115,14 @@ func (bd *BaremetalDesc) IndexKey() string { return bd.Id } +func (bd *BaremetalDesc) AllocCpuNumaPin(vcpuCount, memSizeKB int) []scheduler.SCpuNumaPin { + return nil +} + +func (bd *BaremetalDesc) AllocCpuNumaPinWithNodeCount(vcpuCount, memSizeKB, nodeCount int) []scheduler.SCpuNumaPin { + return nil +} + func (bd *BaremetalDesc) FreeCPUCount() int64 { if bd.ServerID == "" { return int64(bd.CpuCount) diff --git a/pkg/scheduler/cache/candidate/base.go b/pkg/scheduler/cache/candidate/base.go index e68965e52b..9de6ba9790 100644 --- a/pkg/scheduler/cache/candidate/base.go +++ b/pkg/scheduler/cache/candidate/base.go @@ -443,7 +443,7 @@ func (b BaseHostDesc) GetSchedDesc() *jsonutils.JSONDict { func (b *BaseHostDesc) GetPendingUsage() *schedmodels.SPendingUsage { usage, err := schedmodels.HostPendingUsageManager.GetPendingUsage(b.GetId()) if err != nil { - return schedmodels.NewPendingUsageBySchedInfo(b.GetId(), nil) + return schedmodels.NewPendingUsageBySchedInfo(b.GetId(), nil, nil) } return usage } diff --git a/pkg/scheduler/cache/candidate/hosts.go b/pkg/scheduler/cache/candidate/hosts.go index 85342dafa7..aa73e6aa10 100644 --- a/pkg/scheduler/cache/candidate/hosts.go +++ b/pkg/scheduler/cache/candidate/hosts.go @@ -15,10 +15,13 @@ package candidate import ( + "context" "encoding/json" + "sort" gosync "sync" "time" + "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/sets" @@ -26,11 +29,14 @@ import ( "yunion.io/x/sqlchemy" computeapi "yunion.io/x/onecloud/pkg/apis/compute" + hostapi "yunion.io/x/onecloud/pkg/apis/host" + "yunion.io/x/onecloud/pkg/apis/scheduler" computedb "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/compute/baremetal" computemodels "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/scheduler/core" o "yunion.io/x/onecloud/pkg/scheduler/options" + "yunion.io/x/onecloud/pkg/util/cgrouputils/cpuset" ) type hostGetter struct { @@ -65,6 +71,14 @@ func (h *hostGetter) FreeMemorySize(useRsvd bool) int64 { return h.h.GetFreeMemSize(useRsvd) } +func (h *hostGetter) NumaAllocateEnabled() bool { + return h.h.HostTopo.NumaEnabled +} + +func (h *hostGetter) GetFreeCpuNuma() []*scheduler.SFreeNumaCpuMem { + return h.h.GetFreeCpuNuma() +} + func (h *hostGetter) RunningMemorySize() int64 { return h.h.RunningMemSize } @@ -116,6 +130,9 @@ type HostDesc struct { RequiredMemSize int64 `json:"required_mem_size"` FakeDeletedMemSize int64 `json:"fake_deleted_mem_size"` + EnableCpuNumaAllocate bool `json:"enable_cpu_numa_allocate"` + HostTopo *SHostTopo `json:"host_topo"` + // storage StorageTypes []string `json:"storage_types"` @@ -135,6 +152,346 @@ type HostDesc struct { GuestReservedResourceUsed *ReservedResource `json:"guest_reserved_used"` } +type CPUDie struct { + LogicalProcessors cpuset.CPUSet + CpuFree map[int]int + VcpuCount int +} + +func (d *CPUDie) initCpuFree(cpuCmtbound int) { + cpuFree := map[int]int{} + for _, cpuId := range d.LogicalProcessors.ToSliceNoSort() { + cpuFree[cpuId] = cpuCmtbound + } + d.CpuFree = cpuFree +} + +type SorttedCPUDie []*CPUDie + +func (pq SorttedCPUDie) Len() int { return len(pq) } + +func (pq SorttedCPUDie) Less(i, j int) bool { + return pq[i].VcpuCount < pq[j].VcpuCount +} + +func (pq SorttedCPUDie) Swap(i, j int) { + pq[i], pq[j] = pq[j], pq[i] +} + +func (pq *SorttedCPUDie) Push(item interface{}) { + *pq = append(*pq, item.(*CPUDie)) +} + +func (pq *SorttedCPUDie) Pop() interface{} { + old := *pq + n := len(old) + item := old[n-1] + old[n-1] = nil // avoid memory leak + *pq = old[0 : n-1] + return item +} + +func (pq *SorttedCPUDie) LoadCpus(cpus []int, vcpuCount int) { + var cpuDies = map[int][]int{} + for i := 0; i < len(cpus); i++ { + for j := 0; j < len(*pq); j++ { + if (*pq)[j].LogicalProcessors.Contains(cpus[i]) { + if cpuDie, ok := cpuDies[j]; !ok { + cpuDies[j] = []int{cpus[i]} + } else { + cpuDies[j] = append(cpuDie, cpus[i]) + } + break + } + } + } + + for i := 0; i < len(*pq); i++ { + if cpus, ok := cpuDies[i]; ok { + d := (*pq)[i] + for _, cpu := range cpus { + d.CpuFree[cpu] -= 1 + } + d.VcpuCount += vcpuCount + } + } + sort.Sort(pq) +} + +type NumaNode struct { + CpuDies SorttedCPUDie + LogicalProcessors cpuset.CPUSet + VcpuCount int + CpuCount int + + NodeId int + NumaHugeMemSizeKB int + NumaHugeFreeMemSizeKB int +} + +func (n *NumaNode) allocCpuset(vcpuCount int, usedCpu map[int]int) { + for i := range n.CpuDies { + for cpuId, nFree := range n.CpuDies[i].CpuFree { + if cnt, ok := usedCpu[cpuId]; ok { + if cnt < nFree { + usedCpu[cpuId] = cnt + 1 + vcpuCount -= 1 + if vcpuCount <= 0 { + return + } + } + } else { + if nFree > 0 { + usedCpu[cpuId] = 1 + vcpuCount -= 1 + if vcpuCount <= 0 { + return + } + } + } + + } + } + n.allocCpuset(vcpuCount, usedCpu) +} + +func (n *NumaNode) AllocCpuset(vcpuCount int) []int { + var usedCpuCount = make(map[int]int) + n.allocCpuset(vcpuCount, usedCpuCount) + + var ret = make([]int, 0) + for cpuId, cnt := range usedCpuCount { + for cnt > 0 { + ret = append(ret, cpuId) + cnt -= 1 + } + } + return ret +} + +func NewNumaNode(nodeId, hugepageSizeKb int, nodeHugepages []hostapi.HostNodeHugepageNr) *NumaNode { + n := new(NumaNode) + n.LogicalProcessors = cpuset.NewCPUSet() + n.NodeId = nodeId + + for i := range nodeHugepages { + if nodeHugepages[i].NodeId == nodeId { + n.NumaHugeMemSizeKB = nodeHugepages[i].HugepageNr * hugepageSizeKb + } + } + + n.NumaHugeFreeMemSizeKB = n.NumaHugeMemSizeKB + return n +} + +type SHostTopo struct { + Nodes []*NumaNode + NumaEnabled bool + CPUCmtbound int +} + +func (pq SHostTopo) Len() int { return len(pq.Nodes) } + +func (pq SHostTopo) Less(i, j int) bool { + if pq.NumaEnabled { + if pq.Nodes[i].NumaHugeFreeMemSizeKB == pq.Nodes[j].NumaHugeFreeMemSizeKB { + return pq.Nodes[i].VcpuCount < pq.Nodes[j].VcpuCount + } + return pq.Nodes[i].NumaHugeFreeMemSizeKB > pq.Nodes[j].NumaHugeFreeMemSizeKB + } else { + return pq.Nodes[i].VcpuCount < pq.Nodes[j].VcpuCount + } +} + +func (pq SHostTopo) Swap(i, j int) { + pq.Nodes[i], pq.Nodes[j] = pq.Nodes[j], pq.Nodes[i] +} + +func (pq *SHostTopo) Push(item interface{}) { + (*pq).Nodes = append((*pq).Nodes, item.(*NumaNode)) +} + +func (h *SHostTopo) LoadCpuNumaPin(guestsCpuNumaPin []scheduler.SCpuNumaPin) { + for _, gCpuNumaPin := range guestsCpuNumaPin { + var node *NumaNode + for i := range h.Nodes { + if h.Nodes[i].NodeId == gCpuNumaPin.NodeId { + node = h.Nodes[i] + } + } + + cpus := gCpuNumaPin.CpuPin + node.CpuDies.LoadCpus(cpus, len(cpus)) + if h.NumaEnabled && gCpuNumaPin.MemSizeMB != nil { + node.NumaHugeFreeMemSizeKB -= *gCpuNumaPin.MemSizeMB * 1024 + } + node.VcpuCount += len(cpus) + } + sort.Sort(h) +} + +func (h *SHostTopo) nodesEnough(nodeCount, vcpuCount int, memSizeKB int) bool { + var leastFree = memSizeKB / nodeCount + var leastCpuCount = vcpuCount / nodeCount + var remPcpuCount = vcpuCount % nodeCount + + for i := 0; i < nodeCount; i++ { + if h.NumaEnabled { + if h.Nodes[i].NumaHugeFreeMemSizeKB < leastFree { + return false + } + } + + requireCpuCount := leastCpuCount + if remPcpuCount > 0 { + requireCpuCount += 1 + remPcpuCount -= 1 + } + if (h.Nodes[i].VcpuCount + requireCpuCount) > h.Nodes[i].CpuCount*h.CPUCmtbound { + return false + } + + } + return true +} + +func (h *SHostTopo) AllocCpuNumaNodes(vcpuCount, memSizeKB int) []scheduler.SCpuNumaPin { + res := make([]scheduler.SCpuNumaPin, 0) + for nodeCount := 1; nodeCount <= len(h.Nodes); nodeCount *= 2 { + if nodeCount > vcpuCount { + break + } + if ok := h.nodesEnough(nodeCount, vcpuCount, memSizeKB); !ok { + log.Infof("node count %d not enough", nodeCount) + continue + } + log.Infof("use node count %d", nodeCount) + + var nodeAllocSize = memSizeKB / nodeCount + if h.NumaEnabled { + if nodeAllocSize/1024%1024 > 0 { + continue + } + } + + var pcpuCount = vcpuCount / nodeCount + var remPcpuCount = vcpuCount % nodeCount + for i := 0; i < nodeCount; i++ { + var npcpuCount = pcpuCount + if remPcpuCount > 0 { + npcpuCount += 1 + remPcpuCount -= 1 + } + cpuNumaPin := scheduler.SCpuNumaPin{ + CpuPin: h.Nodes[i].AllocCpuset(npcpuCount), + NodeId: h.Nodes[i].NodeId, + } + if h.NumaEnabled { + allocSize := nodeAllocSize / 1024 + cpuNumaPin.MemSizeMB = &allocSize + } + res = append(res, cpuNumaPin) + } + break + } + + return res +} + +func (h *SHostTopo) AllocCpuNumaNodesWithNodeCount(vcpuCount, memSizeKB, nodeCount int) []scheduler.SCpuNumaPin { + res := make([]scheduler.SCpuNumaPin, 0) + var nodeAllocSize = memSizeKB / nodeCount + var pcpuCount = vcpuCount / nodeCount + var remPcpuCount = vcpuCount % nodeCount + for i := 0; i < nodeCount; i++ { + var npcpuCount = pcpuCount + if remPcpuCount > 0 { + npcpuCount += 1 + remPcpuCount -= 1 + } + + cpuNumaPin := scheduler.SCpuNumaPin{ + CpuPin: h.Nodes[i].AllocCpuset(npcpuCount), + NodeId: h.Nodes[i].NodeId, + } + if h.NumaEnabled { + allocSize := nodeAllocSize / 1024 + cpuNumaPin.MemSizeMB = &allocSize + } + res = append(res, cpuNumaPin) + } + return res +} + +func (b *HostBuilder) buildHostTopo( + desc *HostDesc, reservedCpus *cpuset.CPUSet, + hugepageSizeKb int, nodeHugepages []hostapi.HostNodeHugepageNr, + info *hostapi.HostTopology, +) error { + hostTopo := new(SHostTopo) + hostTopo.Nodes = make([]*NumaNode, len(info.Nodes)) + + hasL3Cache := false + for i := 0; i < len(info.Nodes); i++ { + node := NewNumaNode(info.Nodes[i].ID, hugepageSizeKb, nodeHugepages) + + cpuDies := make([]*CPUDie, 0) + for j := 0; j < len(info.Nodes[i].Caches); j++ { + if info.Nodes[i].Caches[j].Level != 3 { + continue + } + hasL3Cache = true + cpuDie := new(CPUDie) + dieBuilder := cpuset.NewBuilder() + for k := 0; k < len(info.Nodes[i].Caches[j].LogicalProcessors); k++ { + if reservedCpus != nil && reservedCpus.Contains(int(info.Nodes[i].Caches[j].LogicalProcessors[k])) { + continue + } + dieBuilder.Add(int(info.Nodes[i].Caches[j].LogicalProcessors[k])) + } + cpuDie.LogicalProcessors = dieBuilder.Result() + cpuDie.initCpuFree(int(desc.CPUCmtbound)) + + node.CpuCount += cpuDie.LogicalProcessors.Size() + node.LogicalProcessors = node.LogicalProcessors.Union(cpuDie.LogicalProcessors) + cpuDies = append(cpuDies, cpuDie) + + // TODO: add cpu core builder + } + if !hasL3Cache { + cpuDie := new(CPUDie) + dieBuilder := cpuset.NewBuilder() + for j := 0; j < len(info.Nodes[i].Cores); j++ { + for k := 0; k < len(info.Nodes[i].Cores[j].LogicalProcessors); k++ { + if reservedCpus != nil && reservedCpus.Contains(info.Nodes[i].Cores[j].LogicalProcessors[k]) { + continue + } + dieBuilder.Add(info.Nodes[i].Cores[j].LogicalProcessors[k]) + } + } + cpuDie.LogicalProcessors = dieBuilder.Result() + node.CpuCount += cpuDie.LogicalProcessors.Size() + node.LogicalProcessors = node.LogicalProcessors.Union(cpuDie.LogicalProcessors) + cpuDies = append(cpuDies, cpuDie) + } + + hasL3Cache = false + node.CpuDies = cpuDies + hostTopo.Nodes[i] = node + } + hostTopo.CPUCmtbound = int(desc.CPUCmtbound) + if len(nodeHugepages) > 0 { + hostTopo.NumaEnabled = true + } + + desc.HostTopo = hostTopo + log.Infof("host topo %s", jsonutils.Marshal(hostTopo)) + + sort.Sort(desc.HostTopo) + desc.EnableCpuNumaAllocate = true + return nil +} + type ReservedResource struct { CPUCount int64 `json:"cpu_count"` MemorySize int64 `json:"memory_size"` @@ -343,6 +700,36 @@ func (h *HostDesc) GetFreeMemSize(useRsvd bool) int64 { return reservedResourceAddCal(h.FreeMemSize, h.GuestReservedMemSizeFree(), useRsvd) - int64(h.GetPendingUsage().Memory) } +func (h *HostDesc) GetFreeCpuNuma() scheduler.SortedFreeNumaCpuMam { + if !h.EnableCpuNumaAllocate { + return nil + } + + res := make(scheduler.SortedFreeNumaCpuMam, 0) + cpuPin := h.GetPendingUsage().CpuPin + numaPin := h.GetPendingUsage().NumaMemPin + for i := range h.HostTopo.Nodes { + nodeFree := new(scheduler.SFreeNumaCpuMem) + nodeFree.NodeId = h.HostTopo.Nodes[i].NodeId + nodeFree.CpuCount = h.HostTopo.Nodes[i].CpuCount + nodeFree.MemSize = h.HostTopo.Nodes[i].NumaHugeFreeMemSizeKB * 1024 + nodeFree.EnableNumaAllocate = h.HostTopo.NumaEnabled + nodeFree.FreeCpuCount = h.HostTopo.Nodes[i].CpuCount*int(h.CPUCmtbound) - h.HostTopo.Nodes[i].VcpuCount + for cpuId, pending := range cpuPin { + if h.HostTopo.Nodes[i].LogicalProcessors.Contains(cpuId) { + nodeFree.FreeCpuCount -= pending + } + } + + if memSize, ok := numaPin[h.HostTopo.Nodes[i].NodeId]; ok { + nodeFree.MemSize -= memSize + } + res = append(res, nodeFree) + } + sort.Sort(res) + return res +} + func (h *HostDesc) GuestReservedMemSizeFree() int64 { return h.GuestReservedResource.MemorySize - h.GuestReservedResourceUsed.MemorySize } @@ -379,6 +766,20 @@ func (h *HostDesc) IndexKey() string { return h.Id } +func (h *HostDesc) AllocCpuNumaPin(vcpuCount, memSizeKB int) []scheduler.SCpuNumaPin { + if !h.EnableCpuNumaAllocate { + return nil + } + return h.HostTopo.AllocCpuNumaNodes(vcpuCount, memSizeKB) +} + +func (h *HostDesc) AllocCpuNumaPinWithNodeCount(vcpuCount, memSizeKB, nodeCount int) []scheduler.SCpuNumaPin { + if !h.EnableCpuNumaAllocate { + return nil + } + return h.HostTopo.AllocCpuNumaNodesWithNodeCount(vcpuCount, memSizeKB, nodeCount) +} + type WaitGroupWrapper struct { gosync.WaitGroup } @@ -694,13 +1095,13 @@ func (b *HostBuilder) InitFuncs() []InitFunc { } } +// build host desc func (b *HostBuilder) BuildOne(host *computemodels.SHost, getter *networkGetter, baseDesc *BaseHostDesc) (interface{}, error) { desc := &HostDesc{ BaseHostDesc: baseDesc, } desc.Metadata = make(map[string]string) - desc.CPUCmtbound = host.GetCPUOvercommitBound() desc.MemCmtbound = host.GetMemoryOvercommitBound() @@ -712,10 +1113,11 @@ func (b *HostBuilder) BuildOne(host *computemodels.SHost, getter *networkGetter, desc.GuestReservedResourceUsed = guestRsvdUsed fillFuncs := []func(*HostDesc, *computemodels.SHost) error{ - b.fillGuestsResourceInfo, //b.fillResidentGroups, b.fillMetadata, b.fillCPUIOLoads, + b.fillGuestsCpuNumaPin, + b.fillGuestsResourceInfo, } for _, f := range fillFuncs { @@ -728,13 +1130,52 @@ func (b *HostBuilder) BuildOne(host *computemodels.SHost, getter *networkGetter, return desc, nil } -func _in(s string, ss []string) bool { - for _, str := range ss { - if s == str { - return true +func (b *HostBuilder) fillGuestsCpuNumaPin(desc *HostDesc, host *computemodels.SHost) error { + if !host.EnableNumaAllocate { + return nil + } + + topoObj, err := host.SysInfo.Get("topology") + if err != nil { + return errors.Wrap(err, "get topology from host sys_info") + } + hostTopo := new(hostapi.HostTopology) + if err := topoObj.Unmarshal(hostTopo); err != nil { + return errors.Wrap(err, "Unmarshal host topology struct") + } + var reservedCpus *cpuset.CPUSet + reservedCpusStr := host.GetMetadata(context.Background(), computeapi.HOSTMETA_RESERVED_CPUS_INFO, nil) + if reservedCpusStr != "" { + reservedCpusJson, err := jsonutils.ParseString(reservedCpusStr) + if err != nil { + return errors.Wrap(err, "parse reserved cpus info failed") + } + reservedCpusInfo := computeapi.HostReserveCpusInput{} + err = reservedCpusJson.Unmarshal(&reservedCpusInfo) + if err != nil { + return errors.Wrap(err, "unmarshal host reserved cpus info failed") + } + reservedCpuset, err := cpuset.Parse(reservedCpusInfo.Cpus) + if err != nil { + return errors.Wrap(err, "cpuset parse reserved cpus") + } + reservedCpus = &reservedCpuset + } + + nodeHugepages := make([]hostapi.HostNodeHugepageNr, 0) + if host.SysInfo.Contains("node_hugepages") { + err = host.SysInfo.Unmarshal(&nodeHugepages, "node_hugepages") + if err != nil { + return errors.Wrap(err, "unmarshal node hugepages") } } - return false + + hugepageSizeKb, err := host.SysInfo.Int("hugepage_size_kb") + if err != nil { + return errors.Wrap(err, "unmarshal hugepage size kb") + } + + return b.buildHostTopo(desc, reservedCpus, int(hugepageSizeKb), nodeHugepages, hostTopo) } func (b *HostBuilder) fillGuestsResourceInfo(desc *HostDesc, host *computemodels.SHost) error { @@ -752,6 +1193,7 @@ func (b *HostBuilder) fillGuestsResourceInfo(desc *HostDesc, host *computemodels creatingMemSize int64 creatingCPUCount int64 creatingGuestCount int64 + guestsCpuNumaPin = make([]scheduler.SCpuNumaPin, 0) ) guestsOnHost, ok := b.hostGuests[host.Id] if !ok { @@ -775,6 +1217,13 @@ func (b *HostBuilder) fillGuestsResourceInfo(desc *HostDesc, host *computemodels runningCount++ memSize += int64(guest.VmemSize) cpuCount += int64(guest.VcpuCount) + if guest.CpuNumaPin != nil { + cpuNumaPin := make([]scheduler.SCpuNumaPin, 0) + if err := guest.CpuNumaPin.Unmarshal(&cpuNumaPin); err != nil { + return errors.Wrap(err, "unmarshal cpu numa pin") + } + guestsCpuNumaPin = append(guestsCpuNumaPin, cpuNumaPin...) + } } else if IsGuestCreating(guest) { creatingGuestCount++ creatingMemSize += int64(guest.VmemSize) @@ -796,6 +1245,11 @@ func (b *HostBuilder) fillGuestsResourceInfo(desc *HostDesc, host *computemodels //} //} } + + if len(guestsCpuNumaPin) > 0 { + desc.HostTopo.LoadCpuNumaPin(guestsCpuNumaPin) + } + desc.GuestCount = guestCount desc.CreatingGuestCount = creatingGuestCount desc.RunningGuestCount = runningCount diff --git a/pkg/scheduler/core/generic_scheduler.go b/pkg/scheduler/core/generic_scheduler.go index 8dad04c55f..4293b23465 100644 --- a/pkg/scheduler/core/generic_scheduler.go +++ b/pkg/scheduler/core/generic_scheduler.go @@ -147,6 +147,7 @@ func (g *GenericScheduler) Schedule(ctx context.Context, unit *Unit, candidates var selectedCandidates []*SelectedCandidate if len(filteredCandidates) > 0 { trace.Step("Prioritizing") + // prioritizing candidates // load all priorities and calculate the candidate's score priorityList, err := PrioritizeCandidates(unit, filteredCandidates, g.priorities) if err != nil { @@ -154,7 +155,7 @@ func (g *GenericScheduler) Schedule(ctx context.Context, unit *Unit, candidates } trace.Step("Selecting hosts") - // select target candate hosts + // select target candidate hosts selectedCandidates, err = SelectHosts(unit, priorityList) if err != nil { return nil, err @@ -603,6 +604,7 @@ func PrioritizeCandidates( results[i] = make(HostPriorityList, len(candidates)) } + // map reduce take priorities processCandidate := func(index int) { var err error candidate := candidates[index] diff --git a/pkg/scheduler/core/result.go b/pkg/scheduler/core/result.go index e7a0c0eac6..166f837512 100644 --- a/pkg/scheduler/core/result.go +++ b/pkg/scheduler/core/result.go @@ -77,13 +77,25 @@ func (its SchedResultItems) Less(i, j int) bool { func (item *SchedResultItem) ToCandidateResource(storageUsed *StorageUsed) *schedapi.CandidateResource { return &schedapi.CandidateResource{ - HostId: item.ID, - Name: item.Name, - Disks: item.getDisks(storageUsed), - Nets: item.Nets, + HostId: item.ID, + CpuNumaPin: item.selectCpuNumaPin(), + Name: item.Name, + Disks: item.getDisks(storageUsed), + Nets: item.Nets, } } +func (item *SchedResultItem) selectCpuNumaPin() []schedapi.SCpuNumaPin { + // if !item.Candidater.Getter().Host().EnableNumaAllocate { + // return nil + // } + + if item.SchedData.LiveMigrate && len(item.SchedData.CpuNumaPin) > 0 { + return item.Candidater.AllocCpuNumaPinWithNodeCount(item.SchedData.Ncpu, item.SchedData.Memory, len(item.SchedData.CpuNumaPin)) + } + return item.Candidater.AllocCpuNumaPin(item.SchedData.Ncpu, item.SchedData.Memory*1024) +} + func (item *SchedResultItem) getDisks(used *StorageUsed) []*schedapi.CandidateDisk { inputs := item.SchedData.Disks ret := make([]*schedapi.CandidateDisk, 0) diff --git a/pkg/scheduler/core/result_helper.go b/pkg/scheduler/core/result_helper.go index 987e4ddc67..15e1bc5d53 100644 --- a/pkg/scheduler/core/result_helper.go +++ b/pkg/scheduler/core/result_helper.go @@ -61,6 +61,7 @@ func transToSchedResult(result *SchedResultItemList, schedInfo *api.SchedInfo) * } } +// trans to region sched results func transToRegionSchedResult(result SchedResultItems, count int64, sid string) *schedapi.ScheduleOutput { apiResults := make([]*schedapi.CandidateResource, 0) succCount := 0 diff --git a/pkg/scheduler/core/types.go b/pkg/scheduler/core/types.go index d1507999a3..0617329f5e 100644 --- a/pkg/scheduler/core/types.go +++ b/pkg/scheduler/core/types.go @@ -97,6 +97,8 @@ type CandidatePropertyGetter interface { RunningMemorySize() int64 TotalMemorySize(useRsvd bool) int64 FreeMemorySize(useRsvd bool) int64 + GetFreeCpuNuma() []*schedapi.SFreeNumaCpuMem + NumaAllocateEnabled() bool StorageInfo() []*baremetal.BaremetalStorage GetFreeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) @@ -140,6 +142,8 @@ type Candidater interface { GetSchedDesc() *jsonutils.JSONDict GetGuestCount() int64 GetResourceType() string + AllocCpuNumaPin(vcpuCount, memSizeKB int) []schedapi.SCpuNumaPin + AllocCpuNumaPinWithNodeCount(vcpuCount, memSizeKB, nodeCount int) []schedapi.SCpuNumaPin } // HostPriority represents the priority of scheduling to particular host, higher priority is better. diff --git a/pkg/scheduler/manager/task_queue.go b/pkg/scheduler/manager/task_queue.go index a48912abbb..350aa7cae6 100644 --- a/pkg/scheduler/manager/task_queue.go +++ b/pkg/scheduler/manager/task_queue.go @@ -84,6 +84,7 @@ func (te *TaskExecutor) Execute(ctx context.Context) { } } +// do execute schedule() func (te *TaskExecutor) execute(ctx context.Context) (*core.ScheduleResult, error) { scheduler := te.scheduler genericScheduler, err := core.NewGenericScheduler(scheduler.(core.Scheduler)) @@ -99,6 +100,7 @@ func (te *TaskExecutor) execute(ctx context.Context) (*core.ScheduleResult, erro te.unit = scheduler.Unit() schedInfo := te.unit.SchedInfo + // generate result helper helper := GenerateResultHelper(schedInfo) result, err := genericScheduler.Schedule(ctx, te.unit, candidates, helper) if err != nil { @@ -108,6 +110,8 @@ func (te *TaskExecutor) execute(ctx context.Context) (*core.ScheduleResult, erro return result, nil } driver := te.unit.GetHypervisorDriver() + + // set sched pending usage if err := setSchedPendingUsage(driver, schedInfo, result.Result); err != nil { return nil, errors.Wrap(err, "setSchedPendingUsage") } diff --git a/pkg/scheduler/models/pending_usage.go b/pkg/scheduler/models/pending_usage.go index 4011f1b4db..03bcf325ab 100644 --- a/pkg/scheduler/models/pending_usage.go +++ b/pkg/scheduler/models/pending_usage.go @@ -50,14 +50,14 @@ func (m *SHostPendingUsageManager) Keyword() string { return "pending_usage_manager" } -func (m *SHostPendingUsageManager) newSessionUsage(req *api.SchedInfo, hostId string) *SessionPendingUsage { +func (m *SHostPendingUsageManager) newSessionUsage(req *api.SchedInfo, hostId string, candidate *schedapi.CandidateResource) *SessionPendingUsage { su := NewSessionUsage(req.SessionId, hostId) - su.Usage = NewPendingUsageBySchedInfo(hostId, req) + su.Usage = NewPendingUsageBySchedInfo(hostId, req, candidate) return su } func (m *SHostPendingUsageManager) newPendingUsage(hostId string) *SPendingUsage { - return NewPendingUsageBySchedInfo(hostId, nil) + return NewPendingUsageBySchedInfo(hostId, nil, nil) } func (m *SHostPendingUsageManager) GetPendingUsage(hostId string) (*SPendingUsage, error) { @@ -81,7 +81,7 @@ func (m *SHostPendingUsageManager) AddPendingUsage(req *api.SchedInfo, candidate sessionUsage, _ := m.GetSessionUsage(req.SessionId, hostId) if sessionUsage == nil { - sessionUsage = m.newSessionUsage(req, hostId) + sessionUsage = m.newSessionUsage(req, hostId, candidate) sessionUsage.StartTimer() } m.addSessionUsage(candidate.HostId, sessionUsage) @@ -100,6 +100,7 @@ func (m *SHostPendingUsageManager) addSessionUsage(hostId string, usage *Session if pendingUsage == nil { pendingUsage = m.newPendingUsage(hostId) } + // add pending usage pendingUsage.Add(usage.Usage) usage.AddCount() m.store.SetSessionUsage(usage.SessionId, hostId, usage) @@ -190,7 +191,7 @@ func NewSessionUsage(sid, hostId string) *SessionPendingUsage { su := &SessionPendingUsage{ HostId: hostId, SessionId: sid, - Usage: NewPendingUsageBySchedInfo(hostId, nil), + Usage: NewPendingUsageBySchedInfo(hostId, nil, nil), count: 0, countLock: new(sync.Mutex), cancelCh: make(chan string), @@ -292,7 +293,9 @@ func (u *SResourcePendingUsage) IsEmpty() bool { type SPendingUsage struct { HostId string Cpu int + CpuPin map[int]int Memory int + NumaMemPin map[int]int IsolatedDevice int DiskUsage *SResourcePendingUsage NetUsage *SResourcePendingUsage @@ -300,7 +303,7 @@ type SPendingUsage struct { InstanceGroupUsage map[string]*api.CandidateGroup } -func NewPendingUsageBySchedInfo(hostId string, req *api.SchedInfo) *SPendingUsage { +func NewPendingUsageBySchedInfo(hostId string, req *api.SchedInfo, candidate *schedapi.CandidateResource) *SPendingUsage { u := &SPendingUsage{ HostId: hostId, DiskUsage: NewResourcePendingUsage(nil), @@ -309,6 +312,8 @@ func NewPendingUsageBySchedInfo(hostId string, req *api.SchedInfo) *SPendingUsag // group init u.InstanceGroupUsage = make(map[string]*api.CandidateGroup) + u.CpuPin = make(map[int]int) + u.NumaMemPin = make(map[int]int) if req == nil { return u @@ -317,6 +322,24 @@ func NewPendingUsageBySchedInfo(hostId string, req *api.SchedInfo) *SPendingUsag u.Memory = req.Memory u.IsolatedDevice = len(req.IsolatedDevices) + if candidate != nil && len(candidate.CpuNumaPin) > 0 { + for _, cpuNumaPin := range candidate.CpuNumaPin { + if cpuNumaPin.MemSizeMB != nil { + if v, ok := u.NumaMemPin[cpuNumaPin.NodeId]; ok { + u.NumaMemPin[cpuNumaPin.NodeId] = v + *cpuNumaPin.MemSizeMB + } + } + + for i := range cpuNumaPin.CpuPin { + if v, ok := u.CpuPin[cpuNumaPin.CpuPin[i]]; ok { + u.CpuPin[cpuNumaPin.CpuPin[i]] = v + 1 + } else { + u.CpuPin[cpuNumaPin.CpuPin[i]] = 1 + } + } + } + } + for _, disk := range req.Disks { backend := disk.Backend size := disk.SizeMb @@ -361,7 +384,22 @@ func (self *SPendingUsage) ToMap() map[string]interface{} { func (self *SPendingUsage) Add(sUsage *SPendingUsage) { self.Cpu = self.Cpu + sUsage.Cpu + for k, v1 := range sUsage.CpuPin { + if v2, ok := self.CpuPin[k]; ok { + self.CpuPin[k] = v1 + v2 + } else { + self.CpuPin[k] = v1 + } + } + self.Memory = self.Memory + sUsage.Memory + for k, v1 := range sUsage.NumaMemPin { + if v2, ok := self.NumaMemPin[k]; ok { + self.NumaMemPin[k] = v1 + v2 + } else { + self.NumaMemPin[k] = v1 + } + } self.IsolatedDevice = self.IsolatedDevice + sUsage.IsolatedDevice self.DiskUsage.Add(sUsage.DiskUsage) self.NetUsage.Add(sUsage.NetUsage) @@ -376,7 +414,19 @@ func (self *SPendingUsage) Add(sUsage *SPendingUsage) { func (self *SPendingUsage) Sub(sUsage *SPendingUsage) { self.Cpu = quotas.NonNegative(self.Cpu - sUsage.Cpu) + for k, v1 := range sUsage.CpuPin { + if v2, ok := self.CpuPin[k]; ok { + self.CpuPin[k] = quotas.NonNegative(v2 - v1) + } + } + self.Memory = quotas.NonNegative(self.Memory - sUsage.Memory) + for k, v1 := range sUsage.NumaMemPin { + if v2, ok := self.NumaMemPin[k]; ok { + self.NumaMemPin[k] = quotas.NonNegative(v2 - v1) + } + } + self.IsolatedDevice = quotas.NonNegative(self.IsolatedDevice - sUsage.IsolatedDevice) self.DiskUsage.Sub(sUsage.DiskUsage) self.NetUsage.Sub(sUsage.NetUsage) diff --git a/pkg/scheduler/test/mock/core.go b/pkg/scheduler/test/mock/core.go index 99a5b6934a..cf95842236 100644 --- a/pkg/scheduler/test/mock/core.go +++ b/pkg/scheduler/test/mock/core.go @@ -34,6 +34,7 @@ import ( jsonutils "yunion.io/x/jsonutils" + "yunion.io/x/onecloud/pkg/apis/scheduler" types "yunion.io/x/onecloud/pkg/cloudcommon/types" baremetal "yunion.io/x/onecloud/pkg/compute/baremetal" models "yunion.io/x/onecloud/pkg/compute/models" @@ -150,6 +151,30 @@ func (mr *MockCandidatePropertyGetterMockRecorder) FreeMemorySize(arg0 interface return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "FreeMemorySize", reflect.TypeOf((*MockCandidatePropertyGetter)(nil).FreeMemorySize), arg0) } +func (m *MockCandidatePropertyGetter) NumaAllocateEnabled() bool { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "NumaAllocateEnabled") + ret0, _ := ret[0].(bool) + return ret0 +} + +func (mr *MockCandidatePropertyGetterMockRecorder) NumaAllocateEnabled() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "NumaAllocateEnabled", reflect.TypeOf((*MockCandidatePropertyGetter)(nil).NumaAllocateEnabled)) +} + +func (m *MockCandidatePropertyGetter) GetFreeCpuNuma() []*scheduler.SFreeNumaCpuMem { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "GetFreeCpuNuma") + ret0, _ := ret[0].([]*scheduler.SFreeNumaCpuMem) + return ret0 +} + +func (mr *MockCandidatePropertyGetterMockRecorder) GetFreeCpuNuma() *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetFreeCpuNuma", reflect.TypeOf((*MockCandidatePropertyGetter)(nil).GetFreeCpuNuma)) +} + // GetFreeGroupCount mocks base method func (m *MockCandidatePropertyGetter) GetFreeGroupCount(arg0 string) (int, error) { m.ctrl.T.Helper() @@ -863,6 +888,30 @@ func (mr *MockCandidaterMockRecorder) IndexKey() *gomock.Call { return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "IndexKey", reflect.TypeOf((*MockCandidater)(nil).IndexKey)) } +func (m *MockCandidater) AllocCpuNumaPin(arg0, arg1 int) []scheduler.SCpuNumaPin { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "AllocCpuNumaPin", arg0, arg1) + ret0, _ := ret[0].([]scheduler.SCpuNumaPin) + return ret0 +} + +func (mr *MockCandidaterMockRecorder) AllocCpuNumaPin(arg0, arg1 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AllocCpuNumaPin", reflect.TypeOf((*MockCandidater)(nil).AllocCpuNumaPin), arg0, arg1) +} + +func (m *MockCandidater) AllocCpuNumaPinWithNodeCount(arg0, arg1, arg2 int) []scheduler.SCpuNumaPin { + m.ctrl.T.Helper() + ret := m.ctrl.Call(m, "AllocCpuNumaPinWithNodeCount", arg0, arg1, arg2) + ret0, _ := ret[0].([]scheduler.SCpuNumaPin) + return ret0 +} + +func (mr *MockCandidaterMockRecorder) AllocCpuNumaPinWithNodeCount(arg0, arg1, arg2 interface{}) *gomock.Call { + mr.mock.ctrl.T.Helper() + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AllocCpuNumaPinWithNodeCount", reflect.TypeOf((*MockCandidater)(nil).AllocCpuNumaPinWithNodeCount), arg0, arg1, arg2) +} + // Type mocks base method func (m *MockCandidater) Type() int { m.ctrl.T.Helper() diff --git a/pkg/scheduler/test/prepare.go b/pkg/scheduler/test/prepare.go index 198e98bb31..e295bdd9f4 100644 --- a/pkg/scheduler/test/prepare.go +++ b/pkg/scheduler/test/prepare.go @@ -83,6 +83,7 @@ func buildCandidate(ctrl *gomock.Controller, param sGetterParams) *mock.MockCand cn.EXPECT().Getter().AnyTimes().Return(getter) cn.EXPECT().IndexKey().AnyTimes().Return(getter.Id()) cn.EXPECT().GetResourceType().AnyTimes().Return(getter.ResourceType()) + cn.EXPECT().AllocCpuNumaPin(gomock.Any(), gomock.Any()).AnyTimes().Return(nil) return cn } @@ -171,6 +172,7 @@ func buildGetter(ctrl *gomock.Controller, param sGetterParams) *mock.MockCandida cg.EXPECT().FreeCPUCount(gomock.Any()).AnyTimes().Return(param.FreeCPUCount) cg.EXPECT().TotalMemorySize(gomock.Any()).AnyTimes().Return(param.TotalMemorySize) cg.EXPECT().FreeMemorySize(gomock.Any()).AnyTimes().Return(param.FreeMemorySize) + cg.EXPECT().GetFreeCpuNuma().AnyTimes().Return(nil) cg.EXPECT().GetFreeStorageSizeOfType(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).AnyTimes().Return(param.FreeStorageSizeAnyType, int64(0), nil) if param.QuotaKeys != nil { cg.EXPECT().GetQuotaKeys(gomock.Any()).AnyTimes().Return(param.QuotaKeys)