diff --git a/pkg/apis/compute/host.go b/pkg/apis/compute/host.go index 9d4b537077..76bf44b652 100644 --- a/pkg/apis/compute/host.go +++ b/pkg/apis/compute/host.go @@ -552,7 +552,8 @@ type SHostPingInput struct { type HostReserveCpusInput struct { Cpus string Mems string - DisableSchedLoadBalance *bool `json:"disable_sched_load_balance"` + DisableSchedLoadBalance *bool `json:"disable_sched_load_balance"` + ProcessesPrefix []string `json:"processes_prefix"` } type HostAutoMigrateInput struct { diff --git a/pkg/apis/compute/host_const.go b/pkg/apis/compute/host_const.go index b6c007bfdb..857b9fd493 100644 --- a/pkg/apis/compute/host_const.go +++ b/pkg/apis/compute/host_const.go @@ -185,4 +185,5 @@ const ( const ( HOSTMETA_RESERVED_CPUS_INFO = "reserved_cpus_info" + HOSTMETA_RESERVED_CPUS_RATE = "reserved_cpus_rate" ) diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index d3dd71bd72..4b7e80b65a 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -4750,6 +4750,50 @@ func (hh *SHost) PerformPing(ctx context.Context, userCred mcclient.TokenCredent return result, nil } +func (host *SHost) getHostNodeReservePercent(reservedCpusStr string) (map[string]float32, error) { + reservedCpuset, err := cpuset.Parse(reservedCpusStr) + if err != nil { + return nil, errors.Wrap(err, "cpuset parse reserved cpus") + } + + topoObj, err := host.SysInfo.Get("topology") + if err != nil { + return nil, errors.Wrap(err, "get topology from host sys_info") + } + info := new(hostapi.HostTopology) + if err := topoObj.Unmarshal(info); err != nil { + return nil, errors.Wrap(err, "Unmarshal host topology struct") + } + nodecpus := map[int]int{} + nodeReservedCpus := map[int]int{} + for i := range info.Nodes { + cSet := 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 reservedCpuset.Contains(info.Nodes[i].Cores[j].LogicalProcessors[k]) { + if cnt, ok := nodeReservedCpus[info.Nodes[i].ID]; !ok { + nodeReservedCpus[info.Nodes[i].ID] = 1 + } else { + nodeReservedCpus[info.Nodes[i].ID] = 1 + cnt + } + } + + cSet.Add(info.Nodes[i].Cores[j].LogicalProcessors[k]) + } + } + nodecpus[info.Nodes[i].ID] = cSet.Result().Size() + } + reserveRate := map[string]float32{} + for nodeId, cnt := range nodecpus { + reserveCnt, ok := nodeReservedCpus[nodeId] + if !ok { + reserveCnt = 0 + } + reserveRate[strconv.Itoa(nodeId)] = float32(reserveCnt) / float32(cnt) + } + return reserveRate, nil +} + func (host *SHost) getHostLogicalCores() ([]int, error) { cpuObj, err := host.SysInfo.Get("cpu_info") if err != nil { @@ -4805,14 +4849,6 @@ func (hh *SHost) PerformReserveCpus( return nil, httperrors.NewNotSupportedError("host type %s not support reserve cpus", hh.HostType) } - cnt, err := hh.GetRunningGuestCount() - if err != nil { - return nil, err - } - if cnt > 0 { - return nil, httperrors.NewBadRequestError("host %s has %d guests, can't update reserve cpus", hh.Id, cnt) - } - if input.Cpus == "" { return nil, httperrors.NewInputParameterError("missing cpus") } @@ -4847,11 +4883,27 @@ func (hh *SHost) PerformReserveCpus( } } + if len(input.Cpus) > 0 { + reservePercent, err := hh.getHostNodeReservePercent(input.Cpus) + if err != nil { + return nil, errors.Errorf("failed getHostNodeReservePercent: %s", err) + } + err = hh.SetMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_RATE, reservePercent, userCred) + if err != nil { + return nil, err + } + } else { + err = hh.RemoveMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_RATE, userCred) + if err != nil { + return nil, err + } + } + err = hh.SetMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_INFO, input, userCred) if err != nil { return nil, err } - if hh.CpuReserved < cs.Size() { + if hh.CpuReserved != cs.Size() { _, err = db.Update(hh, func() error { hh.CpuReserved = cs.Size() return nil diff --git a/pkg/compute/models/isolated_devices.go b/pkg/compute/models/isolated_devices.go index fd75a5b4b1..67fd4562b4 100644 --- a/pkg/compute/models/isolated_devices.go +++ b/pkg/compute/models/isolated_devices.go @@ -21,6 +21,7 @@ import ( "math" "reflect" "sort" + "strconv" "strings" "time" @@ -652,20 +653,18 @@ func (pq *SorttedGroupDevs) Pop() interface{} { return item } -func (manager *SIsolatedDeviceManager) attachHostDeviceToGuestByModel( - ctx context.Context, guest *SGuest, host *SHost, devConfig *api.IsolatedDeviceConfig, - userCred mcclient.TokenCredential, usedDevMap map[string]*SIsolatedDevice, preferNumaNodes []int, -) error { - if len(devConfig.Model) == 0 { - return fmt.Errorf("Not found model from info: %#v", devConfig) - } - // if dev type is not nic, wire is empty string - devs, err := manager.findHostUnusedByDevConfig(devConfig.Model, devConfig.DevType, host.Id, devConfig.WireId) +type SNodeIsolateDevicesInfo struct { + TotalDevCount int + ReservedRate float32 +} + +func (manager *SIsolatedDeviceManager) getDevNodesUsedRate( + ctx context.Context, host *SHost, devConfig *api.IsolatedDeviceConfig, topo *hostapi.HostTopology, +) (map[string]SNodeIsolateDevicesInfo, error) { + devs, err := manager.findHostDevsByDevConfig(devConfig.Model, devConfig.DevType, host.Id, devConfig.WireId) if err != nil || len(devs) == 0 { - return fmt.Errorf("Can't found model %s on host %s", devConfig.Model, host.Id) + return nil, fmt.Errorf("Can't found model %s on host %s", devConfig.Model, host.Id) } - // 1. group devices by device_path - groupDevs := make(SorttedGroupDevs, 0) mapDevs := map[string][]SIsolatedDevice{} for i := range devs { dev := devs[i] @@ -680,11 +679,206 @@ func (manager *SIsolatedDeviceManager) attachHostDeviceToGuestByModel( } mapDevs[devPath] = gdevs } + nodesGroupDevs := map[string]SorttedGroupDevs{} for devPath, mappedDevs := range mapDevs { - groupDevs = append(groupDevs, &GroupDevs{ - DevPath: devPath, - Devs: mappedDevs, - }) + numaNode := strconv.Itoa(int(mappedDevs[0].NumaNode)) + if _, ok := nodesGroupDevs[numaNode]; ok { + nodesGroupDevs[numaNode] = append(nodesGroupDevs[numaNode], &GroupDevs{ + DevPath: devPath, + Devs: mappedDevs, + }) + } else { + groupDevs := make(SorttedGroupDevs, 0) + nodesGroupDevs[numaNode] = append(groupDevs, &GroupDevs{ + DevPath: devPath, + Devs: mappedDevs, + }) + } + } + + reserveRate := map[string]float32{} + reserveRateStr := host.GetMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_RATE, nil) + reserveRateJ, err := jsonutils.ParseString(reserveRateStr) + if err != nil { + return nil, errors.Wrap(err, "parse reserveRateStr") + } + err = reserveRateJ.Unmarshal(&reserveRate) + if err != nil { + return nil, errors.Wrap(err, "unmarshal reserveRateStr") + } + + nodeNoDevIds := map[int]int{} + for i := range topo.Nodes { + nodeId := strconv.Itoa(topo.Nodes[i].ID) + if _, ok := nodesGroupDevs[nodeId]; !ok { + nodeInt, _ := strconv.Atoi(nodeId) + nodeNoDevIds[nodeInt] = -1 + } + } + // + //for nodeId, _ := range reserveRate { + // if _, ok := nodesGroupDevs[nodeId]; !ok { + // nodeInt, _ := strconv.Atoi(nodeId) + // nodeNoDevIds[nodeInt] = -1 + // } + //} + + reserveNodes := map[string][]string{} + for i := range topo.Nodes { + if _, ok := nodeNoDevIds[topo.Nodes[i].ID]; ok { + minDistance := int(math.MaxInt16) + selectNodeId := "" + for nodeId, _ := range nodesGroupDevs { + nodeInt, _ := strconv.Atoi(nodeId) + if topo.Nodes[i].Distances[nodeInt] < minDistance { + selectNodeId = strconv.Itoa(nodeInt) + minDistance = topo.Nodes[i].Distances[nodeInt] + } + } + noDevNodeId := strconv.Itoa(topo.Nodes[i].ID) + log.Debugf("node %s select node %s", noDevNodeId, selectNodeId) + if nodes, ok := reserveNodes[selectNodeId]; ok { + reserveNodes[selectNodeId] = append(nodes, noDevNodeId) + } else { + reserveNodes[selectNodeId] = []string{noDevNodeId} + } + } + } + reserveRates := map[string]SNodeIsolateDevicesInfo{} + for nodeId, devGroups := range nodesGroupDevs { + nodeCnt := 1 + nodeReserveRate := reserveRate[nodeId] + if nodes, ok := reserveNodes[nodeId]; ok { + for i := range nodes { + nodeReserveRate += reserveRate[nodes[i]] + nodeCnt += 1 + } + } + nodeReserveRate = nodeReserveRate / float32(nodeCnt) + devCnt := 0 + for i := range devGroups { + devCnt += len(devGroups[i].Devs) + } + reserveRates[nodeId] = SNodeIsolateDevicesInfo{ + TotalDevCount: devCnt, + ReservedRate: nodeReserveRate, + } + log.Debugf("node %v nodeCnt %v nodeReserveRate %v", nodeId, nodeCnt, nodeReserveRate) + } + return reserveRates, nil +} + +func (manager *SIsolatedDeviceManager) attachHostDeviceToGuestByModel( + ctx context.Context, guest *SGuest, host *SHost, devConfig *api.IsolatedDeviceConfig, + userCred mcclient.TokenCredential, usedDevMap map[string]*SIsolatedDevice, preferNumaNodes []int, +) error { + if len(devConfig.Model) == 0 { + return fmt.Errorf("Not found model from info: %#v", devConfig) + } + // if dev type is not nic, wire is empty string + devs, err := manager.findHostUnusedByDevConfig(devConfig.Model, devConfig.DevType, host.Id, devConfig.WireId) + if err != nil || len(devs) == 0 { + return fmt.Errorf("Can't found model %s on host %s", devConfig.Model, host.Id) + } + // 1. group devices by device_path and numa nodes + //groupDevs := make(SorttedGroupDevs, 0) + mapDevs := map[string][]SIsolatedDevice{} + for i := range devs { + dev := devs[i] + devPath := dev.DevicePath + var gdevs []SIsolatedDevice + + gdevs, ok := mapDevs[devPath] + if !ok { + gdevs = []SIsolatedDevice{dev} + } else { + gdevs = append(gdevs, dev) + } + mapDevs[devPath] = gdevs + } + + var groupDevs SorttedGroupDevs + if len(preferNumaNodes) > 0 { + groupDevs = make(SorttedGroupDevs, 0) + for devPath, mappedDevs := range mapDevs { + groupDevs = append(groupDevs, &GroupDevs{ + DevPath: devPath, + Devs: mappedDevs, + }) + } + } else { + nodesGroupDevs := map[int8]SorttedGroupDevs{} + for devPath, mappedDevs := range mapDevs { + numaNode := mappedDevs[0].NumaNode + if _, ok := nodesGroupDevs[numaNode]; ok { + nodesGroupDevs[numaNode] = append(nodesGroupDevs[numaNode], &GroupDevs{ + DevPath: devPath, + Devs: mappedDevs, + }) + } else { + groupDevs := make(SorttedGroupDevs, 0) + nodesGroupDevs[numaNode] = append(groupDevs, &GroupDevs{ + DevPath: devPath, + Devs: mappedDevs, + }) + } + } + + var selectedNode int8 = -1 + if len(nodesGroupDevs) == 1 { + for nodeId := range nodesGroupDevs { + selectedNode = nodeId + } + } else { + reservedCpusStr := host.GetMetadata(ctx, api.HOSTMETA_RESERVED_CPUS_INFO, nil) + if len(reservedCpusStr) > 0 { + topoObj, err := host.SysInfo.Get("topology") + if err != nil { + return errors.Wrap(err, "get topology from host sys_info") + } + topo := new(hostapi.HostTopology) + if err := topoObj.Unmarshal(topo); err != nil { + return errors.Wrap(err, "Unmarshal host topology struct") + } + nodesReserveRate, err := manager.getDevNodesUsedRate(ctx, host, devConfig, topo) + if err != nil { + return err + } + var selectedNodeUtil float32 = 1.0 + for nodeId, gds := range nodesGroupDevs { + freeDevCnt := 0 + for i := range gds { + freeDevCnt += len(gds[i].Devs) + } + + nodeTotalCnt := nodesReserveRate[strconv.Itoa(int(nodeId))].TotalDevCount + usedDevCnt := nodeTotalCnt - freeDevCnt + + nodeReserveRate := nodesReserveRate[strconv.Itoa(int(nodeId))].ReservedRate + nodeCnt := (1 - nodeReserveRate) * float32(nodeTotalCnt) + nodeutil := float32(usedDevCnt) / nodeCnt + log.Debugf("selectedNodeUtil node %v util %v usedDevCnt %v totalDevCnt %v", nodeId, nodeutil, usedDevCnt, nodeCnt) + if nodeutil < selectedNodeUtil { + selectedNodeUtil = nodeutil + selectedNode = nodeId + } + } + } else { + var selectedNodeDevCnt = 0 + for nodeId, gds := range nodesGroupDevs { + devCnt := 0 + for i := range gds { + devCnt += len(gds[i].Devs) + } + if devCnt > selectedNodeDevCnt { + selectedNodeDevCnt = devCnt + selectedNode = nodeId + } + } + } + } + log.Debugf("selectedNodeUtil node %v", selectedNode) + groupDevs = nodesGroupDevs[selectedNode] } sort.Sort(groupDevs) @@ -845,6 +1039,29 @@ func (manager *SIsolatedDeviceManager) findHostUnusedByDevConfig(model, devType, return manager.findHostUnusedByDevAttr(model, "dev_type", devType, hostId, wireId) } +func (manager *SIsolatedDeviceManager) findHostDevsByDevConfig(model, devType, hostId, wireId string) ([]SIsolatedDevice, error) { + return manager.findHostDevsByDevAttr(model, "dev_type", devType, hostId, wireId) +} +func (manager *SIsolatedDeviceManager) findHostDevsByDevAttr(model, attrKey, attrVal, hostId, wireId string) ([]SIsolatedDevice, error) { + devs := make([]SIsolatedDevice, 0) + q := manager.Query() + q = q.Equals("model", model).Equals("host_id", hostId) + if attrVal != "" { + q.Equals(attrKey, attrVal) + } + if wireId != "" { + wire := WireManager.FetchWireById(wireId) + if wire.VpcId == api.DEFAULT_VPC_ID { + q = q.Equals("wire_id", wireId) + } + } + err := db.FetchModelObjects(manager, q, &devs) + if err != nil { + return nil, err + } + return devs, nil +} + func (manager *SIsolatedDeviceManager) findHostUnusedByDevAttr(model, attrKey, attrVal, hostId, wireId string) ([]SIsolatedDevice, error) { devs := make([]SIsolatedDevice, 0) q := manager.findUnusedQuery() diff --git a/pkg/hostman/guestman/guesthelper.go b/pkg/hostman/guestman/guesthelper.go index 6a72eeabe5..d1ea6150df 100644 --- a/pkg/hostman/guestman/guesthelper.go +++ b/pkg/hostman/guestman/guesthelper.go @@ -219,7 +219,10 @@ type CpuSetCounter struct { Nodes []*NumaNode NumaEnabled bool CPUCmtbound float32 - Lock sync.Mutex + MEMCmtbound float32 + + GuestIds map[string]struct{} + Lock sync.Mutex } func NewGuestCpuSetCounter( @@ -230,6 +233,8 @@ func NewGuestCpuSetCounter( cpuSetCounter.Nodes = make([]*NumaNode, len(info.Nodes)) cpuSetCounter.NumaEnabled = numaAllocate cpuSetCounter.CPUCmtbound = cpuCmtbound + cpuSetCounter.MEMCmtbound = memCmtBound + cpuSetCounter.GuestIds = map[string]struct{}{} hasL3Cache := false nodeReserveMem := reservedMemMb / len(info.Nodes) * 1024 for i := 0; i < len(info.Nodes); i++ { @@ -244,6 +249,7 @@ func NewGuestCpuSetCounter( if err != nil { return nil, err } + reservedCpuCnt := 0 cpuDies := make([]*CPUDie, 0) for j := 0; j < len(info.Nodes[i].Caches); j++ { if info.Nodes[i].Caches[j].Level != 3 { @@ -254,6 +260,7 @@ func NewGuestCpuSetCounter( 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])) { + reservedCpuCnt += 1 continue } dieBuilder.Add(int(info.Nodes[i].Caches[j].LogicalProcessors[k])) @@ -271,6 +278,7 @@ func NewGuestCpuSetCounter( 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]) { + reservedCpuCnt += 1 continue } dieBuilder.Add(info.Nodes[i].Cores[j].LogicalProcessors[k]) @@ -286,6 +294,7 @@ func NewGuestCpuSetCounter( hasL3Cache = false node.CpuDies = cpuDies + node.ReserveCpuCount = reservedCpuCnt sort.Sort(node.CpuDies) cpuSetCounter.Nodes[i] = node } @@ -294,13 +303,14 @@ func NewGuestCpuSetCounter( return cpuSetCounter, nil } -func (pq *CpuSetCounter) AllocCpusetWithNodeCount(vcpuCount int, memSizeKB int64, nodeCount int) (map[int]SAllocNumaCpus, error) { +func (pq *CpuSetCounter) AllocCpusetWithNodeCount(vcpuCount int, memSizeKB int64, nodeCount int, guestId string) (map[int]SAllocNumaCpus, error) { if !pq.NumaEnabled { - return pq.AllocCpuset(vcpuCount, memSizeKB, nil) + return pq.AllocCpuset(vcpuCount, memSizeKB, nil, guestId) } if len(pq.Nodes) < nodeCount { return nil, nil } + pq.GuestIds[guestId] = struct{}{} pq.Lock.Lock() defer pq.Lock.Unlock() @@ -339,13 +349,14 @@ func (pq *CpuSetCounter) IsNumaEnabled() bool { return pq.NumaEnabled } -func (pq *CpuSetCounter) AllocCpuset(vcpuCount int, memSizeKB int64, preferNumaNodes []int8) (map[int]SAllocNumaCpus, error) { +func (pq *CpuSetCounter) AllocCpuset(vcpuCount int, memSizeKB int64, preferNumaNodes []int8, guestId string) (map[int]SAllocNumaCpus, error) { pq.Lock.Lock() defer pq.Lock.Unlock() if len(pq.Nodes) == 0 { return nil, nil } + pq.GuestIds[guestId] = struct{}{} if pq.NumaEnabled && len(preferNumaNodes) > 0 { sortedNumaDistance := pq.getDistancesSeqByPreferNodes(preferNumaNodes) @@ -416,7 +427,7 @@ func (pq *CpuSetCounter) allocCpuNumaNodesByPreferNodes( } else { log.Infof("node %v not enough", pq.Nodes[i]) } - log.Infof("node %d, free mems %d", pq.Nodes[nodeIdx].NodeId, pq.Nodes[nodeIdx].NumaNodeFreeMemSizeKB) + log.Infof("node %d, free mems %d, vcpuCount %d, GuestCounts %v", pq.Nodes[nodeIdx].NodeId, pq.Nodes[nodeIdx].NumaNodeFreeMemSizeKB, pq.Nodes[nodeIdx].VcpuCount, len(pq.GuestIds)) } if allocatedNode < nodeCount { @@ -429,6 +440,8 @@ type SSortedNumaDistance struct { NodeIndex int Distance int FreeMemSize int + UsedRate float32 + CpuReserved bool } func (pq *CpuSetCounter) getDistancesSeqByPreferNodes(preferNumaNodes []int8) []SSortedNumaDistance { @@ -438,19 +451,37 @@ func (pq *CpuSetCounter) getDistancesSeqByPreferNodes(preferNumaNodes []int8) [] for j := range preferNumaNodes { distance += pq.Nodes[i].Distances[preferNumaNodes[j]] } + + var useableCpuRate float32 = 1.0 + if pq.Nodes[i].ReserveCpuCount > 0 { + useableCpuRate = float32(pq.Nodes[i].CpuCount) / float32(pq.Nodes[i].CpuCount+pq.Nodes[i].ReserveCpuCount) + } + + usedMems := float32(pq.Nodes[i].NumaNodeMemSizeKB - pq.Nodes[i].NumaNodeFreeMemSizeKB) + usedRate := usedMems / (float32(pq.Nodes[i].MemTotalSizeKB) * pq.MEMCmtbound * useableCpuRate) + + //memCmt := float32(usedMems / pq.Nodes[i].NumaNodeMemSizeKB) + //cpuPro := float32(pq.Nodes[i].CpuCount) * pq.CPUCmtbound / (float32(pq.Nodes[i].CpuCount)*pq.CPUCmtbound - float32(pq.Nodes[i].VcpuCount)) sortedNumaDistance[i] = SSortedNumaDistance{ NodeIndex: i, Distance: distance, FreeMemSize: int(pq.Nodes[i].NumaNodeFreeMemSizeKB), + UsedRate: usedRate, + CpuReserved: pq.Nodes[i].ReserveCpuCount > 0, } } sort.Slice(sortedNumaDistance, func(i, j int) bool { // 7 is tolerant max distances - if (sortedNumaDistance[i].Distance + 7) < sortedNumaDistance[j].Distance { + if sortedNumaDistance[i].Distance > (7 + sortedNumaDistance[j].Distance) { + return false + } else if (sortedNumaDistance[i].Distance + 7) < sortedNumaDistance[j].Distance { return true - } else { - return sortedNumaDistance[i].FreeMemSize > sortedNumaDistance[j].FreeMemSize } + + if sortedNumaDistance[i].CpuReserved { + return sortedNumaDistance[i].UsedRate < sortedNumaDistance[j].UsedRate + } + return sortedNumaDistance[i].FreeMemSize > sortedNumaDistance[j].FreeMemSize }) return sortedNumaDistance } @@ -670,10 +701,12 @@ type NumaNode struct { LogicalProcessors cpuset.CPUSet VcpuCount int CpuCount int + ReserveCpuCount int NodeId int Distances []int NumaNodeMemSizeKB int64 + MemTotalSizeKB int64 NumaNodeFreeMemSizeKB int64 } @@ -696,6 +729,7 @@ func NewNumaNode( return nil, errors.Errorf("node %d no memory info: %#v", nodeInfo.ID, nodeInfo) } n.NumaNodeMemSizeKB = int64(float32(nodeInfo.Memory.TotalUsableBytes/1024-int64(reservedMemSizeKB)) * memCmtBound) + n.MemTotalSizeKB = nodeInfo.Memory.TotalUsableBytes / 1024 } else { nodeHugepagePath := fmt.Sprintf("/sys/devices/system/node/node%d/hugepages/hugepages-%dkB", n.NodeId, hugepageSizeKB) if !fileutils2.Exists(nodeHugepagePath) { diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index 26f434aee3..dd1716e66e 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -2488,7 +2488,7 @@ func (task *SGuestHotplugCpuMemTask) startAddCpusWithFreeVcpuSet(vcpuSet []int) } } } else { - cpus, _ := task.manager.cpuSet.AllocCpuset(1, 0, nil) + cpus, _ := task.manager.cpuSet.AllocCpuset(1, 0, nil, task.GetId()) for _, cpus := range cpus { //pcpus := cpuset.NewCPUSet(cpus.Cpuset...).String() //vcpus := fmt.Sprintf("%d-%d", vcpuId, vcpuId) diff --git a/pkg/hostman/guestman/pod.go b/pkg/hostman/guestman/pod.go index 5c4b040e3d..6d8d862535 100644 --- a/pkg/hostman/guestman/pod.go +++ b/pkg/hostman/guestman/pod.go @@ -1117,7 +1117,7 @@ func (s *sPodGuestInstance) allocateCpuNumaPin() error { } } - nodeNumaCpus, err := s.manager.cpuSet.AllocCpuset(int(s.Desc.Cpu), s.Desc.Mem*1024, preferNumaNodes) + nodeNumaCpus, err := s.manager.cpuSet.AllocCpuset(int(s.Desc.Cpu), s.Desc.Mem*1024, preferNumaNodes, s.GetId()) if err != nil { return err } @@ -1140,6 +1140,12 @@ func (s *sPodGuestInstance) allocateCpuNumaPin() error { vcpuPin := make([]desc.SVCpuPin, len(numaCpus.Cpuset)) for i := range numaCpus.Cpuset { vcpuPin[i].Pcpu = numaCpus.Cpuset[i] + if i < int(s.Desc.Cpu) { + vcpuPin[i].Vcpu = i + } else { + vcpuPin[i].Vcpu = -1 + } + } memPin := &desc.SCpuNumaPin{ diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index 19f3e6c63b..b8ae2017ad 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -177,7 +177,7 @@ func (s *SKVMGuestInstance) reallocateNumaNodes(isMigrate bool) error { } func (s *SKVMGuestInstance) reallocateMigrateNumaNodes() error { - nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(s.Desc.Cpu), s.Desc.Mem*1024, len(s.Desc.MemDesc.Mem.Mems)+1) + nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(s.Desc.Cpu), s.Desc.Mem*1024, len(s.Desc.MemDesc.Mem.Mems)+1, s.GetId()) if err != nil { return errors.Wrap(err, "AllocCpusetWithNodeCount") } @@ -371,7 +371,7 @@ func (s *SKVMGuestInstance) initLiveDescFromSourceGuest(srcDesc *desc.SGuestDesc cpuNumaPin = s.Desc.CpuNumaPin } else { // allocate cpu numa pin local - nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(srcDesc.Cpu), srcDesc.Mem*1024, len(srcDesc.MemDesc.Mem.Mems)+1) + nodeNumaCpus, err := s.manager.cpuSet.AllocCpusetWithNodeCount(int(srcDesc.Cpu), srcDesc.Mem*1024, len(srcDesc.MemDesc.Mem.Mems)+1, s.GetId()) if err != nil { return errors.Wrap(err, "AllocCpusetWithNodeCount") } @@ -2700,7 +2700,7 @@ func (s *SKVMGuestInstance) allocGuestNumaCpuset() error { } } - nodeNumaCpus, err := s.manager.cpuSet.AllocCpuset(int(s.Desc.Cpu), s.Desc.Mem*1024, preferNumaNodes) + nodeNumaCpus, err := s.manager.cpuSet.AllocCpuset(int(s.Desc.Cpu), s.Desc.Mem*1024, preferNumaNodes, s.GetId()) if err != nil { return err } diff --git a/pkg/hostman/guestman/runtime.go b/pkg/hostman/guestman/runtime.go index 924fdf4080..b813f930b6 100644 --- a/pkg/hostman/guestman/runtime.go +++ b/pkg/hostman/guestman/runtime.go @@ -173,6 +173,8 @@ func LoadGuestCpuset(m *SGuestManager, s GuestRuntimeInstance) error { if s.IsRunning() { m.cpuSet.Lock.Lock() defer m.cpuSet.Lock.Unlock() + m.cpuSet.GuestIds[s.GetId()] = struct{}{} + for _, vcpuPin := range guestDesc.VcpuPin { pcpuSet, err := cpuset.Parse(vcpuPin.Pcpus) if err != nil { @@ -186,13 +188,18 @@ func LoadGuestCpuset(m *SGuestManager, s GuestRuntimeInstance) error { } m.cpuSet.LoadCpus(pcpuSet.ToSlice(), vcpuSet.Size()) } + for _, numaCpuset := range guestDesc.CpuNumaPin { pcpus := make([]int, 0) + vcpus := make([]int, 0) for i := range numaCpuset.VcpuPin { pcpus = append(pcpus, numaCpuset.VcpuPin[i].Pcpu) + if numaCpuset.VcpuPin[i].Vcpu >= 0 { + vcpus = append(vcpus, numaCpuset.VcpuPin[i].Vcpu) + } } - m.cpuSet.LoadNumaCpus(numaCpuset.SizeMB, int(*numaCpuset.NodeId), pcpus, len(numaCpuset.VcpuPin)) + m.cpuSet.LoadNumaCpus(numaCpuset.SizeMB, int(*numaCpuset.NodeId), pcpus, len(vcpus)) } } return nil @@ -205,16 +212,22 @@ func ReleaseCpuNumaPin(m *SGuestManager, cpuNumaPin []*desc.SCpuNumaPin) { for _, numaCpus := range cpuNumaPin { pcpus := make([]int, 0) + vcpus := make([]int, 0) for i := range numaCpus.VcpuPin { pcpus = append(pcpus, numaCpus.VcpuPin[i].Pcpu) + if numaCpus.VcpuPin[i].Vcpu >= 0 { + vcpus = append(vcpus, numaCpus.VcpuPin[i].Vcpu) + } } - m.cpuSet.ReleaseNumaCpus(numaCpus.SizeMB, int(*numaCpus.NodeId), pcpus, len(numaCpus.VcpuPin)) + m.cpuSet.ReleaseNumaCpus(numaCpus.SizeMB, int(*numaCpus.NodeId), pcpus, len(vcpus)) } } func ReleaseGuestCpuset(m *SGuestManager, s GuestRuntimeInstance) { m.cpuSet.Lock.Lock() defer m.cpuSet.Lock.Unlock() + delete(m.cpuSet.GuestIds, s.GetId()) + guestDesc := s.GetDesc() ReleaseCpuNumaPin(m, guestDesc.CpuNumaPin) diff --git a/pkg/hostman/hostinfo/hostinfo.go b/pkg/hostman/hostinfo/hostinfo.go index 670f2125ce..022d5418e9 100644 --- a/pkg/hostman/hostinfo/hostinfo.go +++ b/pkg/hostman/hostinfo/hostinfo.go @@ -18,6 +18,7 @@ import ( "bytes" "context" "fmt" + "io/ioutil" "math" "net" "os" @@ -833,6 +834,9 @@ func (h *SHostInfo) initCgroup() error { !reservedCpusTask.CustomConfig(cgrouputils.CPUSET_SCHED_LOAD_BALANCE, "0") { return fmt.Errorf("failed init host reserved cpuset sched load balance") } + if len(h.reservedCpusInfo.ProcessesPrefix) > 0 { + go h.startBindReservedCpus(h.reservedCpusInfo.ProcessesPrefix) + } } return nil } @@ -2670,6 +2674,58 @@ func (h *SHostInfo) MemCmtBound() float32 { return h.memCmtBound } +func (h *SHostInfo) getProcessesPids(processesPrefix []string) (map[string]string, error) { + files, err := ioutil.ReadDir("/proc") + if err != nil { + return nil, err + } + res := map[string]string{} + re := regexp.MustCompile(`^\d+$`) + for _, f := range files { + if re.MatchString(f.Name()) { + cmdline, err := fileutils2.FileGetContents(path.Join("/proc", f.Name(), "cmdline")) + if err != nil { + log.Errorf("failed read proc %s cmdline: %s", f.Name(), err) + continue + } + segs := strings.Split(cmdline, "\x00") + if utils.IsInStringArray(segs[0], processesPrefix) { + res[segs[0]] = f.Name() + log.Infof("getProcessesPids append %s %s", segs[0], f.Name()) + } + } + } + return res, nil +} + +func (h *SHostInfo) startBindReservedCpus(processesPrefix []string) { + for { + processPids, err := h.getProcessesPids(processesPrefix) + if err != nil { + log.Errorf("getProcessesPids %s", err) + } else { + for process, pid := range processPids { + cgroupName := path.Join(hostconsts.HOST_RESERVED_CPUSET, strings.ReplaceAll(process, "/", "_")) + task := cgrouputils.NewCGroupCPUSetTask(pid, cgroupName, 0, "") + if !task.Configure() { + log.Errorf("process failed init reserved cpuset %s %s", process, pid) + continue + } + + if !task.CustomConfig(cgrouputils.CPUSET_CLONE_CHILDREN, "1") { + log.Errorf("process failed set host reserved cpuset clone children %s %s", process, pid) + continue + } + if !task.SetTask() { + log.Errorf("process %s %s failed set cgroup cpuset", process, pid) + continue + } + } + } + time.Sleep(time.Second * 100) + } +} + func NewHostInfo() (*SHostInfo, error) { var res = new(SHostInfo) res.sysinfo = &SSysInfo{} diff --git a/pkg/mcclient/options/compute/host.go b/pkg/mcclient/options/compute/host.go index 0aa751051d..938affdce3 100644 --- a/pkg/mcclient/options/compute/host.go +++ b/pkg/mcclient/options/compute/host.go @@ -92,10 +92,11 @@ type HostReserveCpusOptions struct { Cpus string Mems string DisableSchedLoadBalance bool + ProcessesPrefix []string `help:"Processes prefix bind reserved cpus"` } func (o *HostReserveCpusOptions) Params() (jsonutils.JSONObject, error) { - return options.StructToParams(o) + return jsonutils.Marshal(o), nil } type HostAutoMigrateOnHostDownOptions struct { diff --git a/pkg/scheduler/cache/candidate/hosts.go b/pkg/scheduler/cache/candidate/hosts.go index a8fe31ced5..b550d730ce 100644 --- a/pkg/scheduler/cache/candidate/hosts.go +++ b/pkg/scheduler/cache/candidate/hosts.go @@ -612,11 +612,13 @@ func (h *SHostTopo) getDistancesSeqByPreferNodes(preferNumaNodes []int) []SSorte } sort.Slice(sortedNumaDistance, func(i, j int) bool { // 7 is tolerant max distances - if (sortedNumaDistance[i].Distance + 7) < sortedNumaDistance[j].Distance { + if sortedNumaDistance[i].Distance > (7 + sortedNumaDistance[j].Distance) { + return false + } else if (sortedNumaDistance[i].Distance + 7) < sortedNumaDistance[j].Distance { return true - } else { - return sortedNumaDistance[i].FreeMemSize > sortedNumaDistance[j].FreeMemSize } + + return sortedNumaDistance[i].FreeMemSize > sortedNumaDistance[j].FreeMemSize }) return sortedNumaDistance } diff --git a/pkg/util/cgrouputils/cgrouputils.go b/pkg/util/cgrouputils/cgrouputils.go index bd132fed3b..f672d71aec 100644 --- a/pkg/util/cgrouputils/cgrouputils.go +++ b/pkg/util/cgrouputils/cgrouputils.go @@ -625,6 +625,7 @@ const ( CPUSET_CPUS = "cpuset.cpus" CPUSET_MEMS = "cpuset.mems" CPUSET_SCHED_LOAD_BALANCE = "cpuset.sched_load_balance" + CPUSET_CLONE_CHILDREN = "cgroup.clone_children" ) func (c *CGroupCPUSetTask) Module() string {