feat(host): cgroupv2 support (#23304)

This commit is contained in:
wanyaoqi
2025-09-16 23:20:18 +08:00
committed by GitHub
parent 104c927e74
commit 7a82d04f70
17 changed files with 1023 additions and 371 deletions
+10 -32
View File
@@ -2237,7 +2237,7 @@ func (s *SKVMGuestInstance) ExitCleanup(clear bool) {
}
func (s *SKVMGuestInstance) CleanupCpuset() {
cgPath := path.Join(cgrouputils.RootTaskPath("cpuset"), s.GetCgroupName())
cgPath := path.Join(cgrouputils.GetSubModulePath("cpuset"), s.GetCgroupName())
cgName := s.GetCgroupName()
cgFiles, err := ioutil.ReadDir(cgPath)
if err != nil {
@@ -2248,13 +2248,13 @@ func (s *SKVMGuestInstance) CleanupCpuset() {
continue
}
subCgName := path.Join(cgName, fi.Name())
task := cgrouputils.NewCGroupCPUSetTask(strconv.Itoa(s.GetPid()), subCgName, 0, "")
task := cgrouputils.NewCGroupCPUSetTask(strconv.Itoa(s.GetPid()), subCgName, "", "")
if !task.RemoveTask() {
log.Warningf("remove cpuset cgroup error: %s, pid: %d", s.Id, s.GetPid())
}
}
task := cgrouputils.NewCGroupCPUSetTask(strconv.Itoa(s.GetPid()), s.GetCgroupName(), 0, "")
task := cgrouputils.NewCGroupCPUSetTask(strconv.Itoa(s.GetPid()), s.GetCgroupName(), "", "")
if !task.RemoveTask() {
log.Warningf("remove cpuset cgroup error: %s, pid: %d", s.Id, s.GetPid())
}
@@ -2865,7 +2865,7 @@ func (s *SKVMGuestInstance) GetCgroupName() string {
func (s *SKVMGuestInstance) GuestPrelaunchSetCgroup() {
s.cgroupPid = s.GetPid()
s.setCgroupIo()
//s.setCgroupIo()
s.setCgroupCpu()
if err := s.setCgroupCPUSet(); err != nil {
log.Errorf("Guest %s failed set cgroup cpuset: %s", s.GetName(), err)
@@ -2876,32 +2876,13 @@ func (s *SKVMGuestInstance) setCgroupPid() {
s.cgroupPid = s.GetPid()
}
func (s *SKVMGuestInstance) setCgroupIo() {
appTags := s.getApptags()
params := map[string]int{}
if utils.IsInStringArray("io_hardlimit", appTags) {
devId := s.getStorageDeviceId()
if len(devId) == 0 {
log.Errorln("failed to get device ID (MAJOR:MINOR)")
return
}
params["blkio.throttle.read_bps_device"] = options.HostOptions.DefaultReadBpsPerCpu
params["blkio.throttle.read_iops_device"] = options.HostOptions.DefaultReadIopsPerCpu
params["blkio.throttle.write_bps_device"] = options.HostOptions.DefaultWriteBpsPerCpu
params["blkio.throttle.write_iops_device"] = options.HostOptions.DefaultWriteIopsPerCpu
cgrouputils.CgroupIoHardlimitSet(
strconv.Itoa(s.cgroupPid), s.GetCgroupName(), int(s.Desc.Cpu), params, devId,
)
}
}
func (s *SKVMGuestInstance) setCgroupCpu() {
var (
cpu = s.Desc.Cpu
cpuWeight = 1024
)
cgrouputils.CgroupSet(strconv.Itoa(s.cgroupPid), s.GetCgroupName(), int(cpu)*cpuWeight)
cgrouputils.NewCGroupCPUTask(strconv.Itoa(s.cgroupPid), s.GetCgroupName(), int(cpu)*cpuWeight).SetTask()
}
func (s *SKVMGuestInstance) setCgroupCPUSet() error {
@@ -2932,7 +2913,7 @@ func (s *SKVMGuestInstance) setCgroupCPUSet() error {
guestPid := strconv.Itoa(s.GetPid())
// guest root cpuset group
task := cgrouputils.NewCGroupCPUSetTask(guestPid, cgName, 0, cpusetStr)
task := cgrouputils.NewCGroupCPUSetTask(guestPid, cgName, cpusetStr, "")
if !task.SetTask() {
return errors.Errorf("Cgroup cpuset task failed")
}
@@ -2954,7 +2935,8 @@ func (s *SKVMGuestInstance) setCgroupCPUSet() error {
}
pcpu := s.Desc.CpuNumaPin[i].VcpuPin[j].Pcpu
vcpuCgname := path.Join(cgName, vcpuThreadId)
taskVcpu := cgrouputils.NewCGroupSubCPUSetTask(guestPid, vcpuCgname, 0, strconv.Itoa(pcpu), []string{vcpuThreadId})
taskVcpu := cgrouputils.NewCGroupSubCPUSetTask(guestPid, vcpuCgname, strconv.Itoa(pcpu), []string{vcpuThreadId})
if !taskVcpu.SetTask() {
return errors.Errorf("Vcpu set cgroup cpuset task failed")
}
@@ -3753,9 +3735,7 @@ func (s *SKVMGuestInstance) CPUSet(ctx context.Context, input []int) (*api.Serve
cpusetStr = strings.Join(cpus, ",")
}
task := cgrouputils.NewCGroupCPUSetTask(
strconv.Itoa(s.GetPid()), s.GetCgroupName(), 0, cpusetStr,
)
task := cgrouputils.NewCGroupCPUSetTask(strconv.Itoa(s.GetPid()), s.GetCgroupName(), cpusetStr, "")
if !task.SetTask() {
return nil, errors.Errorf("Cgroup cpuset task failed")
}
@@ -3806,9 +3786,7 @@ func (s *SKVMGuestInstance) CPUSetRemove(ctx context.Context) error {
if !s.IsRunning() {
return nil
}
task := cgrouputils.NewCGroupCPUSetTask(
strconv.Itoa(s.GetPid()), s.GetCgroupName(), 0, "",
)
task := cgrouputils.NewCGroupCPUSetTask(strconv.Itoa(s.GetPid()), s.GetCgroupName(), "", "")
if !task.RemoveTask() {
return errors.Errorf("Remove task error happened, please lookup host log")
}
+6 -11
View File
@@ -442,9 +442,10 @@ func (h *SHostInfo) prepareEnv() error {
log.Warningf("modprobe vhost_net error: %s", output)
}
if !options.HostOptions.DisableSetCgroup {
if !cgrouputils.Init(h.IoScheduler) {
return fmt.Errorf("Cannot initialize control group subsystem")
if err := cgrouputils.Init(h.IoScheduler); err != nil {
return fmt.Errorf("Cannot initialize control group subsystem: %s", err)
}
h.sysinfo.CgroupVersion = cgrouputils.GetCgroupVersion()
}
// err = h.resetIptables()
@@ -768,23 +769,17 @@ func (h *SHostInfo) initCgroup() error {
hostCpuset := hostCpusetBuilder.Result()
hostCpusetStr := hostCpuset.String()
// init host cpuset root group
if !cgrouputils.NewCGroupCPUSetTask("", hostconsts.HOST_CGROUP, 0, hostCpusetStr).Configure() {
if !cgrouputils.NewCGroupCPUSetTask("", hostconsts.HOST_CGROUP, hostCpusetStr, "").Configure() {
return fmt.Errorf("failed init host root cpuset")
}
// init host cpu root group
cgrouputils.CgroupSet("", hostconsts.HOST_CGROUP, hostCpuset.Size()*1024)
// init host blkio root group
cgrouputils.CgroupIoHardlimitSet("", hostconsts.HOST_CGROUP, 0, nil, "")
cgrouputils.NewCGroupCPUTask("", hostconsts.HOST_CGROUP, hostCpuset.Size()*1024).SetTask()
if h.reservedCpusInfo != nil {
reservedCpusTask := cgrouputils.NewCGroupCPUSetTask("", hostconsts.HOST_RESERVED_CPUSET, 0, h.reservedCpusInfo.Cpus)
reservedCpusTask := cgrouputils.NewCGroupCPUSetTask("", hostconsts.HOST_RESERVED_CPUSET, h.reservedCpusInfo.Cpus, h.reservedCpusInfo.Mems)
if !reservedCpusTask.Configure() {
return fmt.Errorf("failed init host reserved cpuset %s", h.reservedCpusInfo.Cpus)
}
if h.reservedCpusInfo.Mems != "" &&
!reservedCpusTask.CustomConfig(cgrouputils.CPUSET_MEMS, h.reservedCpusInfo.Mems) {
return fmt.Errorf("failed init host reserved cpuset mems %s", h.reservedCpusInfo.Mems)
}
if h.reservedCpusInfo.DisableSchedLoadBalance != nil &&
*h.reservedCpusInfo.DisableSchedLoadBalance &&
!reservedCpusTask.CustomConfig(cgrouputils.CPUSET_SCHED_LOAD_BALANCE, "0") {
+1
View File
@@ -329,6 +329,7 @@ type SSysInfo struct {
KvmModule string `json:"kvm_module"`
CpuModelName string `json:"cpu_model_name"`
CpuMicrocode string `json:"cpu_microcode"`
CgroupVersion string `json:"cgroup_version"`
StorageType string `json:"storage_type"`
+5 -13
View File
@@ -30,7 +30,7 @@ import (
schedapi "yunion.io/x/onecloud/pkg/apis/scheduler"
"yunion.io/x/onecloud/pkg/cloudcommon/cmdline"
"yunion.io/x/onecloud/pkg/mcclient/options"
"yunion.io/x/onecloud/pkg/util/cgrouputils"
"yunion.io/x/onecloud/pkg/util/cgrouputils/cpuset"
)
var ErrEmtptyUpdate = errors.New("No valid update data")
@@ -1470,20 +1470,12 @@ type ServerCPUSetOptions struct {
}
func (o *ServerCPUSetOptions) Params() (jsonutils.JSONObject, error) {
sets := cgrouputils.ParseCpusetStr(o.SETS)
parts := strings.Split(sets, ",")
if len(parts) == 0 {
return nil, errors.New(fmt.Sprintf("Invalid cpu sets %q", o.SETS))
cpus, err := cpuset.Parse(o.SETS)
if err != nil {
return nil, errors.New(fmt.Sprintf("parse cpuset failed: %s", err))
}
input := &computeapi.ServerCPUSetInput{
CPUS: make([]int, 0),
}
for _, s := range parts {
sd, err := strconv.Atoi(s)
if err != nil {
return nil, errors.New(fmt.Sprintf("Not digit part %q", s))
}
input.CPUS = append(input.CPUS, sd)
CPUS: cpus.ToSlice(),
}
return jsonutils.Marshal(input), nil
}
+1
View File
@@ -0,0 +1 @@
package cgroup // import "yunion.io/x/onecloud/pkg/util/cgrouputils/cgroup"
+40
View File
@@ -0,0 +1,40 @@
// 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 cgroup
const (
CGROUP_V1 = "cgroup_v1"
CGROUP_V2 = "cgroup_v2"
)
type ICGroupTask interface {
InitTask(hand ICGroupTask, cpuShares int, pid, name string)
SetPid(string)
SetName(string)
SetWeight(coreNum int)
SetHand(hand ICGroupTask)
GetParam(name string) string
CustomConfig(key, value string) bool
GetStaticConfig() map[string]string
GetConfig() map[string]string
Module() string
RemoveTask() bool
SetTask() bool
Configure() bool
TaskIsExist() bool
Init() bool
}
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package cgrouputils
package cgroupv1
import (
"bufio"
@@ -26,184 +26,105 @@ import (
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/util/cgrouputils/cgroup"
"yunion.io/x/onecloud/pkg/util/fileutils2"
"yunion.io/x/onecloud/pkg/util/procutils"
)
const (
CGROUP_PATH_SYSFS = "/sys/fs/cgroup"
CGROUP_PATH_ROOT = "/cgroup"
CGROUP_TASKS = "tasks"
CGROUP_TASKS = "tasks"
maxWeight = 16
normalizeBase = 1024
)
var (
cgroupsPath = getGroupPath()
)
type ICGroupTask interface {
InitTask(hand ICGroupTask, coreNum int, pid, name string)
SetPid(string)
SetName(string)
SetWeight(coreNum int)
SetHand(hand ICGroupTask)
GetStaticConfig() map[string]string
GetConfig() map[string]string
Module() string
RemoveTask() bool
SetTask() bool
Configure() bool
init() bool
}
type CGroupTask struct {
pid string
threadIds []string
name string
weight float64
hand ICGroupTask
hand cgroup.ICGroupTask
}
func NewCGroupTask(pid, name string, coreNum int, threadIds []string) *CGroupTask {
func NewCGroupTask(pid, name string, cpuShares int, threadIds []string) *CGroupTask {
return &CGroupTask{
pid: pid,
name: name,
weight: float64(coreNum) / normalizeBase,
weight: float64(cpuShares) / normalizeBase,
threadIds: threadIds,
}
}
func getGroupPath() string {
if fileutils2.Exists(CGROUP_PATH_SYSFS) {
return CGROUP_PATH_SYSFS
} else {
return CGROUP_PATH_ROOT
}
}
func CgroupIsMounted() bool {
return procutils.NewCommand("mountpoint", cgroupsPath).Run() == nil
}
func ModuleIsMounted(module string) bool {
fullPath := path.Join(cgroupsPath, module)
if fi, err := os.Lstat(fullPath); err != nil {
log.Errorln(err)
return false
} else if fi.Mode()&os.ModeSymlink == os.ModeSymlink {
// is link
fullPath, err = filepath.EvalSymlinks(fullPath)
if err != nil {
log.Errorln(err)
}
}
return procutils.NewCommand("mountpoint", fullPath).Run() == nil
}
func RootTaskPath(module string) string {
return path.Join(cgroupsPath, module)
}
func GetTaskParamPath(module, name, pid string) string {
spath := RootTaskPath(module)
if len(pid) > 0 {
spath = path.Join(spath, pid)
}
return path.Join(spath, name)
}
func GetRootParam(module, name, pid string) string {
param, err := fileutils2.FileGetContents(GetTaskParamPath(module, name, pid))
if err != nil {
log.Errorln(err)
return ""
}
return strings.TrimSpace(param)
}
// cpuset, task, tid, cgname
func SetRootParam(module, name, value, pid string) bool {
param := GetRootParam(module, name, pid)
if param != value {
err := ioutil.WriteFile(GetTaskParamPath(module, name, pid), []byte(value), 0644)
if err != nil {
if len(pid) == 0 {
pid = "root"
func (*CGroupTask) Init() bool {
if !manager.CgroupIsMounted() {
if !fileutils2.Exists(manager.GetCgroupPath()) {
if err := procutils.NewCommand("mkdir", "-p", manager.GetCgroupPath()).Run(); err != nil {
log.Errorf("mkdir -p %s error: %v", manager.GetCgroupPath(), err)
}
log.Errorf("fail to set %s to %s(%s): %s", name, value, pid, err)
}
if err := procutils.NewCommand("mount", "-t", "tmpfs", "-o", "uid=0,gid=0,mode=0755",
"cgroup", manager.GetCgroupPath()).Run(); err != nil {
log.Errorf("mount cgroups path %s, error: %v", manager.GetCgroupPath(), err)
return false
}
}
return true
}
// cleanup
func CleanupNonexistPids(module string, subName string) {
var root = RootTaskPath(module)
cleanNonexitPidsWithRoot(root)
if subName != "" {
root = path.Join(RootTaskPath(module), subName)
cleanNonexitPidsWithRoot(root)
}
}
func cleanNonexitPidsWithRoot(root string) {
files, err := ioutil.ReadDir(root)
file, err := os.Open("/proc/cgroups")
if err != nil {
log.Errorf("GetTaskIds failed: %s", err)
return
}
ids := []string{}
for _, file := range files {
ids = append(ids, file.Name())
log.Errorln(err)
return false
}
defer file.Close()
re1 := regexp.MustCompile(`^\d+$`)
re2 := regexp.MustCompile(`^server_[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}_\d+$`)
for _, pid := range ids {
spid := ""
if re1.MatchString(pid) {
spid = pid
} else if re2.MatchString(pid) {
segs := strings.Split(pid, "_")
spid = segs[len(segs)-1]
}
if fileutils2.IsDir(path.Join(root, pid)) {
if !fileutils2.Exists(path.Join("/proc", spid)) {
log.Infof("Cgroup clenup %s", pid)
subFiles, err := ioutil.ReadDir(path.Join(root, pid))
if err != nil {
log.Errorf("sub dir %s GetTaskIds failed: %s", path.Join(root, pid), err)
} else {
for _, fi := range subFiles {
if !fi.IsDir() {
continue
}
if err := os.Remove(path.Join(root, pid, fi.Name())); err != nil {
log.Errorf("CleanupNonexistPids pid=%s tid=%s error: %s", pid, fi.Name(), err)
}
re := regexp.MustCompile(`\s+`)
scanner := bufio.NewScanner(file)
for scanner.Scan() {
// #subsys_name hierarchy num_cgroups enabled
line := scanner.Text()
if line[0] != '#' {
parts := re.Split(line, -1)
module := parts[0]
if parts[3] == "0" { // disabled
continue
}
if !manager.ModuleIsMounted(module) {
moduleDir := path.Join(manager.GetCgroupPath(), module)
if !fileutils2.Exists(moduleDir) {
if output, err := procutils.NewCommand("mkdir", moduleDir).Output(); err != nil {
log.Errorf("mkdir %s failed: %s, %s", moduleDir, err, output)
return false
}
}
if err := os.Remove(path.Join(root, pid)); err != nil {
log.Errorf("CleanupNonexistPids pid=%s error: %s", pid, err)
if err := procutils.NewCommand("mount", "-t", "cgroup", "-o",
module, module, moduleDir).Run(); err != nil {
log.Errorf("mount cgroup module %s to %s error: %v", module, moduleDir, err)
return false
}
}
}
}
if err := scanner.Err(); err != nil {
log.Errorf("scan file %s error: %v", file.Name(), err)
return false
}
return true
}
func (c *CGroupTask) CustomConfig(key, value string) bool {
configPath := GetTaskParamPath(c.hand.Module(), key, c.GroupName())
if !fileutils2.Exists(configPath) {
return true
}
return SetRootParam(c.hand.Module(), key, value, c.GroupName())
}
func (c *CGroupTask) SetWeight(coreNum int) {
c.weight = float64(coreNum) / normalizeBase
}
func (c *CGroupTask) SetHand(hand ICGroupTask) {
func (c *CGroupTask) SetHand(hand cgroup.ICGroupTask) {
c.hand = hand
}
@@ -222,7 +143,7 @@ func (c *CGroupTask) GroupName() string {
return c.pid
}
func (c *CGroupTask) InitTask(hand ICGroupTask, coreNum int, pid, name string) {
func (c *CGroupTask) InitTask(hand cgroup.ICGroupTask, coreNum int, pid, name string) {
c.SetHand(hand)
c.SetWeight(coreNum)
c.SetPid(pid)
@@ -266,7 +187,7 @@ func (c *CGroupTask) TaskPath() string {
return path.Join(RootTaskPath(c.hand.Module()), c.GroupName())
}
func (c *CGroupTask) taskIsExist() bool {
func (c *CGroupTask) TaskIsExist() bool {
return fileutils2.Exists(c.TaskPath())
}
@@ -295,7 +216,7 @@ func (c *CGroupTask) MoveTasksToRoot() {
}
func (c *CGroupTask) RemoveTask() bool {
if c.taskIsExist() {
if c.TaskIsExist() {
c.MoveTasksToRoot()
log.Infof("Remove task path %s %s", c.TaskPath(), c.name)
if err := os.Remove(c.TaskPath()); err != nil {
@@ -329,7 +250,7 @@ func (c *CGroupTask) SetParams(conf map[string]string) bool {
}
func (c *CGroupTask) Configure() bool {
if !c.taskIsExist() {
if !c.TaskIsExist() {
if !c.createTask() {
return false
}
@@ -345,7 +266,7 @@ func (c *CGroupTask) Configure() bool {
}
func (c *CGroupTask) SetTask() bool {
if !c.taskIsExist() {
if !c.TaskIsExist() {
if !c.createTask() {
return false
}
@@ -396,62 +317,6 @@ func (c *CGroupTask) PushPid(tid string, isRoot bool) {
}
}
func (c *CGroupTask) init() bool {
if !CgroupIsMounted() {
if !fileutils2.Exists(cgroupsPath) {
if err := procutils.NewCommand("mkdir", "-p", cgroupsPath).Run(); err != nil {
log.Errorf("mkdir -p %s error: %v", cgroupsPath, err)
}
}
if err := procutils.NewCommand("mount", "-t", "tmpfs", "-o", "uid=0,gid=0,mode=0755",
"cgroup", cgroupsPath).Run(); err != nil {
log.Errorf("mount cgroups path %s, error: %v", cgroupsPath, err)
return false
}
}
file, err := os.Open("/proc/cgroups")
if err != nil {
log.Errorln(err)
return false
}
defer file.Close()
re := regexp.MustCompile(`\s+`)
scanner := bufio.NewScanner(file)
for scanner.Scan() {
// #subsys_name hierarchy num_cgroups enabled
line := scanner.Text()
if line[0] != '#' {
parts := re.Split(line, -1)
module := parts[0]
if parts[3] == "0" { // disabled
continue
}
if !ModuleIsMounted(module) {
moduleDir := path.Join(cgroupsPath, module)
if !fileutils2.Exists(moduleDir) {
if output, err := procutils.NewCommand("mkdir", moduleDir).Output(); err != nil {
log.Errorf("mkdir %s failed: %s, %s", moduleDir, err, output)
return false
}
}
if err := procutils.NewCommand("mount", "-t", "cgroup", "-o",
module, module, moduleDir).Run(); err != nil {
log.Errorf("mount cgroup module %s to %s error: %v", module, moduleDir, err)
return false
}
}
}
}
if err := scanner.Err(); err != nil {
log.Errorf("scan file %s error: %v", file.Name(), err)
return false
}
return true
}
/**
* CGroupCPUTask
*/
@@ -474,15 +339,15 @@ func (c *CGroupCPUTask) GetConfig() map[string]string {
return map[string]string{CPU_SHARES: fmt.Sprintf("%d", wt)}
}
func (c *CGroupCPUTask) init() bool {
func (c *CGroupCPUTask) Init() bool {
return SetRootParam(c.Module(), CPU_SHARES,
fmt.Sprintf("%d", CgroupsSharesWeight), "")
}
func NewCGroupCPUTask(pid, name string, coreNum int) CGroupCPUTask {
cgroup := CGroupCPUTask{NewCGroupTask(pid, name, coreNum, nil)}
cgroup.hand = &cgroup
return cgroup
func (m *cgroupManager) NewCGroupCPUTask(pid, name string, cpuShares int) cgroup.ICGroupTask {
t := &CGroupCPUTask{NewCGroupTask(pid, name, cpuShares, nil)}
t.SetHand(t)
return t
}
/**
@@ -493,13 +358,6 @@ type CGroupIOTask struct {
*CGroupTask
}
var (
/*
* A global IoScheduler variable, should be set before initialize CGROUP
*/
IoScheduler string
)
const (
IoWeightBase = 100
IoWeightMax = 1000
@@ -523,7 +381,7 @@ func (c *CGroupIOTask) GetConfig() map[string]string {
} else if wt < IoWeightMin {
wt = IoWeightMin
}
switch IoScheduler {
switch manager.GetIoScheduler() {
case IOSCHED_CFQ:
return map[string]string{BLOCK_IO_WEIGHT: fmt.Sprintf("%d", wt)}
case IOSCHED_BFQ:
@@ -533,8 +391,8 @@ func (c *CGroupIOTask) GetConfig() map[string]string {
}
}
func (c *CGroupIOTask) init() bool {
switch IoScheduler {
func (c *CGroupIOTask) Init() bool {
switch manager.GetIoScheduler() {
case IOSCHED_CFQ:
return SetRootParam(c.Module(), BLOCK_IO_WEIGHT, fmt.Sprintf("%d", IoWeightMax), "")
default:
@@ -542,8 +400,8 @@ func (c *CGroupIOTask) init() bool {
}
}
func NewCGroupIOTask(pid, name string, coreNum int) *CGroupIOTask {
task := &CGroupIOTask{NewCGroupTask(pid, name, coreNum, nil)}
func (m *cgroupManager) NewCGroupIOTask(pid, name string, cpuShares int) cgroup.ICGroupTask {
task := &CGroupIOTask{NewCGroupTask(pid, name, cpuShares, nil)}
task.SetHand(task)
return task
}
@@ -570,9 +428,9 @@ func (c *CGroupIOHardlimitTask) GetConfig() map[string]string {
return config
}
func NewCGroupIOHardlimitTask(pid, name string, coreNum int, params map[string]int, devId string) *CGroupIOHardlimitTask {
func (m *cgroupManager) NewCGroupIOHardlimitTask(pid, name string, coreNum int, params map[string]int, devId string) cgroup.ICGroupTask {
task := &CGroupIOHardlimitTask{
CGroupIOTask: NewCGroupIOTask(pid, name, 0),
CGroupIOTask: m.NewCGroupIOTask(pid, name, 0).(*CGroupIOTask),
cpuNum: coreNum,
params: params,
devId: devId,
@@ -603,7 +461,7 @@ func (c *CGroupMemoryTask) GetConfig() map[string]string {
return map[string]string{MEMORY_SWAPPINESS: fmt.Sprintf("%d", vm_swappiness)}
}
func NewCGroupMemoryTask(pid, name string, coreNum int) *CGroupMemoryTask {
func (m *cgroupManager) NewCGroupMemoryTask(pid, name string, coreNum int) cgroup.ICGroupTask {
task := &CGroupMemoryTask{
CGroupTask: NewCGroupTask(pid, name, coreNum, nil),
}
@@ -619,112 +477,45 @@ type CGroupCPUSetTask struct {
*CGroupTask
cpuset string
mems string
}
const (
CPUSET_CPUS = "cpuset.cpus"
CPUSET_MEMS = "cpuset.mems"
CPUSET_SCHED_LOAD_BALANCE = "cpuset.sched_load_balance"
CPUSET_CPUS = "cpuset.cpus"
CPUSET_MEMS = "cpuset.mems"
)
func (c *CGroupCPUSetTask) Module() string {
return "cpuset"
}
func (c *CGroupCPUSetTask) GetStaticConfig() map[string]string {
return map[string]string{CPUSET_MEMS: GetRootParam(c.Module(), CPUSET_MEMS, "")}
}
func (c *CGroupCPUSetTask) GetConfig() map[string]string {
if c.cpuset == "" {
parentPath := filepath.Dir(c.GroupName())
c.cpuset = GetRootParam(c.Module(), CPUSET_CPUS, parentPath)
}
return map[string]string{CPUSET_CPUS: c.cpuset}
}
func (c *CGroupCPUSetTask) CustomConfig(key, value string) bool {
return SetRootParam(c.hand.Module(), key, value, c.GroupName())
}
func NewCGroupCPUSetTask(pid, name string, coreNum int, cpuset string) CGroupCPUSetTask {
task := CGroupCPUSetTask{
CGroupTask: NewCGroupTask(pid, name, coreNum, nil),
cpuset: cpuset,
if c.mems == "" {
parentPath := filepath.Dir(c.GroupName())
c.mems = GetRootParam(c.Module(), CPUSET_MEMS, parentPath)
}
task.SetHand(&task)
return map[string]string{CPUSET_CPUS: c.cpuset, CPUSET_MEMS: c.mems}
}
func (m *cgroupManager) NewCGroupCPUSetTask(pid, name, cpuset, mems string) cgroup.ICGroupTask {
task := &CGroupCPUSetTask{
CGroupTask: NewCGroupTask(pid, name, 0, nil),
cpuset: cpuset,
mems: mems,
}
task.SetHand(task)
return task
}
func NewCGroupSubCPUSetTask(pid, name string, coreNum int, cpuset string, threadIds []string) CGroupCPUSetTask {
task := CGroupCPUSetTask{
CGroupTask: NewCGroupTask(pid, name, coreNum, threadIds),
func (m *cgroupManager) NewCGroupSubCPUSetTask(pid, name string, cpuset string, threadIds []string) cgroup.ICGroupTask {
task := &CGroupCPUSetTask{
CGroupTask: NewCGroupTask(pid, name, 0, threadIds),
cpuset: cpuset,
}
task.SetHand(&task)
task.SetHand(task)
return task
}
func Init(ioScheduler string) bool {
IoScheduler = ioScheduler
for _, hand := range []ICGroupTask{&CGroupTask{}, &CGroupCPUTask{}, &CGroupIOTask{}} {
if !hand.init() {
return false
}
}
return true
}
func CgroupSet(pid, name string, coreNum int) bool {
tasks := []ICGroupTask{
&CGroupCPUTask{&CGroupTask{}},
&CGroupIOTask{&CGroupTask{}},
&CGroupMemoryTask{&CGroupTask{}},
}
for _, hand := range tasks {
hand.InitTask(hand, coreNum, pid, name)
if !hand.SetTask() {
return false
}
}
return true
}
func CgroupIoHardlimitSet(
pid, name string, coreNum int,
params map[string]int, devId string,
) bool {
cg := NewCGroupIOHardlimitTask(pid, name, coreNum, params, devId)
return cg.SetTask()
}
func CgroupDestroy(pid, name string) bool {
tasks := []ICGroupTask{
&CGroupCPUTask{&CGroupTask{}},
&CGroupIOTask{&CGroupTask{}},
&CGroupMemoryTask{&CGroupTask{}},
//&CGroupCPUSetTask{&CGroupTask{}, ""},
&CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}},
}
for _, hand := range tasks {
hand.InitTask(hand, 0, pid, name)
if !hand.RemoveTask() {
return false
}
}
return true
}
func CgroupCleanAll(subName string) {
tasks := []ICGroupTask{
&CGroupCPUTask{&CGroupTask{}},
&CGroupIOTask{&CGroupTask{}},
&CGroupMemoryTask{&CGroupTask{}},
&CGroupCPUSetTask{CGroupTask: &CGroupTask{}},
&CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}},
}
for _, hand := range tasks {
hand.SetHand(hand)
CleanupNonexistPids(hand.Module(), subName)
}
}
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package cgrouputils
package cgroupv1
import (
"bufio"
@@ -28,6 +28,7 @@ func TestCgroupSet(t *testing.T) {
pid, _ := reader.ReadString('\n')
pid = strings.TrimSpace(pid)
t.Logf("Start %s cgroup set", pid)
CgroupSet(pid, "", 1)
CgroupCleanAll("")
Init("/sys/fs/cgroup", "")
manager.NewCGroupCPUTask(pid, "", 1).SetTask()
manager.CgroupCleanAll("")
}
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package cgrouputils
package cgroupv1
import (
"encoding/json"
@@ -38,7 +38,7 @@ var (
utilHistory map[string][]float64
)
func RebalanceProcesses(pids []string) {
func (m *cgroupManager) RebalanceProcesses(pids []string) {
rebalanceProcessesLock.Lock()
if rebalanceProcessesRunning {
rebalanceProcessesLock.Unlock()
@@ -95,7 +95,7 @@ func CommitProcessCpuset(proc *ProcessCPUinfo, idx int) {
cpu, _ := GetSystemCpu()
sets := cpu.GetCpuset(idx)
if len(sets) > 0 {
cpuset := NewCGroupCPUSetTask(strconv.Itoa(proc.Pid), "", 0, sets)
cpuset := manager.NewCGroupCPUSetTask(strconv.Itoa(proc.Pid), "", sets, "")
cpuset.SetTask()
}
}
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package cgrouputils
package cgroupv1
import (
"fmt"
@@ -231,8 +231,8 @@ func NewProcessCPUinfo(pid int) (*ProcessCPUinfo, error) {
cpuinfo.Pid = pid
spid := strconv.Itoa(pid)
cpuTask := NewCGroupCPUTask(spid, "", 0)
if cpuTask.taskIsExist() {
cpuTask := manager.NewCGroupCPUTask(spid, "", 0)
if cpuTask.TaskIsExist() {
share := cpuTask.GetParam("cpu.shares")
ishare, err := strconv.ParseFloat(share, 64)
if err != nil {
@@ -243,8 +243,8 @@ func NewProcessCPUinfo(pid int) (*ProcessCPUinfo, error) {
}
}
cpusetTask := NewCGroupCPUSetTask(fmt.Sprintf("%d", pid), "", 0, "")
if cpusetTask.taskIsExist() {
cpusetTask := manager.NewCGroupCPUSetTask(fmt.Sprintf("%d", pid), "", "", "")
if cpusetTask.TaskIsExist() {
cpuset := cpusetTask.GetParam("cpuset.cpus")
if len(cpuset) > 0 {
c, err := GetSystemCpu()
@@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.
package cgrouputils
package cgroupv1
import (
"os"
+15
View File
@@ -0,0 +1,15 @@
// 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 cgroupv1 // import "yunion.io/x/onecloud/pkg/util/cgrouputils/cgroupv1"
+207
View File
@@ -0,0 +1,207 @@
package cgroupv1
import (
"io/ioutil"
"os"
"path"
"path/filepath"
"regexp"
"strings"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/util/cgrouputils/cgroup"
"yunion.io/x/onecloud/pkg/util/fileutils2"
"yunion.io/x/onecloud/pkg/util/procutils"
)
var (
manager *cgroupManager
)
type cgroupManager struct {
cgroupPath string
ioScheduler string
}
func (m *cgroupManager) GetCgroupPath() string {
return m.cgroupPath
}
func (m *cgroupManager) GetSubModulePath(module string) string {
return path.Join(m.cgroupPath, module)
}
func (m *cgroupManager) GetCgroupVersion() string {
return cgroup.CGROUP_V1
}
func (m *cgroupManager) GetIoScheduler() string {
return m.ioScheduler
}
func (m *cgroupManager) CgroupIsMounted() bool {
return procutils.NewCommand("mountpoint", m.cgroupPath).Run() == nil
}
func (m *cgroupManager) ModuleIsMounted(module string) bool {
fullPath := path.Join(m.cgroupPath, module)
if fi, err := os.Lstat(fullPath); err != nil {
log.Errorln(err)
return false
} else if fi.Mode()&os.ModeSymlink == os.ModeSymlink {
// is link
fullPath, err = filepath.EvalSymlinks(fullPath)
if err != nil {
log.Errorln(err)
}
}
return procutils.NewCommand("mountpoint", fullPath).Run() == nil
}
func Init(cgroupPath, ioScheduler string) (*cgroupManager, error) {
if manager != nil {
return manager, nil
}
manager = &cgroupManager{
cgroupPath: cgroupPath,
ioScheduler: ioScheduler,
}
for _, hand := range []cgroup.ICGroupTask{
&CGroupTask{},
&CGroupCPUTask{},
//&CGroupIOTask{},
} {
if !hand.Init() {
return manager, errors.Errorf("Cannot initialize %s control group subsystem", hand.Module())
}
}
return manager, nil
}
func RootTaskPath(module string) string {
return path.Join(manager.cgroupPath, module)
}
func GetTaskParamPath(module, name, pid string) string {
spath := RootTaskPath(module)
if len(pid) > 0 {
spath = path.Join(spath, pid)
}
return path.Join(spath, name)
}
func GetRootParam(module, name, pid string) string {
param, err := fileutils2.FileGetContents(GetTaskParamPath(module, name, pid))
if err != nil {
log.Errorln(err)
return ""
}
return strings.TrimSpace(param)
}
// cpuset, task, tid, cgname
func SetRootParam(module, name, value, pid string) bool {
param := GetRootParam(module, name, pid)
if param != value {
err := ioutil.WriteFile(GetTaskParamPath(module, name, pid), []byte(value), 0644)
if err != nil {
if len(pid) == 0 {
pid = "root"
}
log.Errorf("fail to set %s to %s(%s): %s", name, value, pid, err)
return false
}
}
return true
}
func (m *cgroupManager) CgroupDestroy(pid, name string) bool {
tasks := []cgroup.ICGroupTask{
&CGroupCPUTask{&CGroupTask{}},
&CGroupIOTask{&CGroupTask{}},
&CGroupMemoryTask{&CGroupTask{}},
//&CGroupCPUSetTask{&CGroupTask{}, ""},
&CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}},
}
for _, hand := range tasks {
hand.InitTask(hand, 0, pid, name)
if !hand.RemoveTask() {
return false
}
}
return true
}
func (m *cgroupManager) CgroupCleanAll(subName string) {
tasks := []cgroup.ICGroupTask{
&CGroupCPUTask{&CGroupTask{}},
&CGroupIOTask{&CGroupTask{}},
&CGroupMemoryTask{&CGroupTask{}},
&CGroupCPUSetTask{CGroupTask: &CGroupTask{}},
&CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}},
}
for _, hand := range tasks {
hand.SetHand(hand)
cleanupNonexistPids(hand.Module(), subName)
}
}
// cleanup
func cleanupNonexistPids(module string, subName string) {
var root = RootTaskPath(module)
cleanNonexitPidsWithRoot(root)
if subName != "" {
root = path.Join(RootTaskPath(module), subName)
cleanNonexitPidsWithRoot(root)
}
}
func cleanNonexitPidsWithRoot(root string) {
files, err := ioutil.ReadDir(root)
if err != nil {
log.Errorf("GetTaskIds failed: %s", err)
return
}
ids := []string{}
for _, file := range files {
ids = append(ids, file.Name())
}
re1 := regexp.MustCompile(`^\d+$`)
re2 := regexp.MustCompile(`^server_[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}_\d+$`)
for _, pid := range ids {
spid := ""
if re1.MatchString(pid) {
spid = pid
} else if re2.MatchString(pid) {
segs := strings.Split(pid, "_")
spid = segs[len(segs)-1]
}
if fileutils2.IsDir(path.Join(root, pid)) {
if !fileutils2.Exists(path.Join("/proc", spid)) {
log.Infof("Cgroup cleanup %s", path.Join(root, pid))
subFiles, err := ioutil.ReadDir(path.Join(root, pid))
if err != nil {
log.Errorf("sub dir %s GetTaskIds failed: %s", path.Join(root, pid), err)
} else {
for _, fi := range subFiles {
if !fi.IsDir() {
continue
}
if err := os.Remove(path.Join(root, pid, fi.Name())); err != nil {
log.Errorf("CleanupNonexistPids pid=%s tid=%s error: %s", pid, fi.Name(), err)
}
}
}
if err := os.Remove(path.Join(root, pid)); err != nil {
log.Errorf("CleanupNonexistPids pid=%s error: %s", pid, err)
}
}
}
}
}
+346
View File
@@ -0,0 +1,346 @@
// 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 cgroupv2
import (
"fmt"
"os"
"path"
"path/filepath"
"regexp"
"strings"
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/util/cgrouputils/cgroup"
"yunion.io/x/onecloud/pkg/util/fileutils2"
)
const (
CGROUP_SUBTREE_CONTROL = "cgroup.subtree_control"
CGROUP_PROCS = "cgroup.procs"
CGROUP_THREADS = "cgroup.threads"
CGROUP_TYPE = "cgroup.type"
CGROUP_TYPE_THREADED = "threaded"
CPU_WEIGHT = "cpu.weight"
)
type CgroupTask struct {
pid string
threadIds []string
name string
hand cgroup.ICGroupTask
}
func NewCGroupBaseTask(pid, name string, threadIds []string) *CgroupTask {
return &CgroupTask{
pid: pid,
name: name,
threadIds: threadIds,
}
}
func (c *CgroupTask) InitTask(hand cgroup.ICGroupTask, cpuShares int, pid, name string) {}
func (c *CgroupTask) SetPid(string) {}
func (c *CgroupTask) SetName(string) {}
func (c *CgroupTask) SetWeight(coreNum int) {}
func (c *CgroupTask) Init() bool {
//initSubGroups()
return true
}
func (c *CgroupTask) Module() string {
return ""
}
func (c *CgroupTask) SetHand(hand cgroup.ICGroupTask) {
c.hand = hand
}
func (c *CgroupTask) GetStaticConfig() map[string]string {
return nil
}
func (c *CgroupTask) GetConfig() map[string]string {
return nil
}
func (c *CgroupTask) CustomConfig(key, value string) bool {
configPath := path.Join(manager.GetCgroupPath(), c.GroupName(), key)
if !fileutils2.Exists(configPath) {
return true
}
return c.SetParam(key, value)
}
func (c *CgroupTask) GetParentGroup() string {
return filepath.Dir(c.GroupName())
}
func (c *CgroupTask) GetParam(name string) string {
return getParam(name, c.GroupName())
}
func (c *CgroupTask) RemoveTask() bool {
if c.TaskIsExist() {
if c.GetParam(CGROUP_TYPE) == CGROUP_TYPE_THREADED {
// move threads to parents
threads := c.GetParam(CGROUP_THREADS)
if len(threads) > 0 {
parentGroup := c.GetParentGroup()
for _, thread := range strings.Split(threads, "\n") {
thread = strings.TrimSpace(thread)
if len(thread) > 0 {
err := setParam(CGROUP_THREADS, thread, parentGroup)
if err != nil {
log.Errorf("failed remove thread %s of %s: %s", thread, c.GroupName(), err)
}
}
}
}
} else {
// move procs to root
procs := c.GetParam(CGROUP_PROCS)
if len(procs) > 0 {
for _, proc := range strings.Split(procs, "\n") {
proc = strings.TrimSpace(proc)
if len(proc) > 0 {
err := setParam(CGROUP_PROCS, proc)
if err != nil {
log.Errorf("failed remove proc %s of %s: %s", proc, c.GroupName(), err)
}
}
}
}
}
log.Infof("Remove task path %s %s", c.TaskPath(), c.name)
if err := os.Remove(c.TaskPath()); err != nil {
log.Errorf("Remove task path failed %s", err)
return false
}
}
return true
}
func (c *CgroupTask) Configure() bool {
if !c.ensureTask() {
return false
}
conf := c.hand.GetStaticConfig()
if !c.SetParams(conf) {
return false
}
conf = c.hand.GetConfig()
return c.SetParams(conf)
}
func (c *CgroupTask) SetTask() bool {
if !c.ensureTask() {
return false
}
conf := c.hand.GetStaticConfig()
if len(conf) > 0 {
if !c.SetParams(conf) {
return false
}
}
conf = c.hand.GetConfig()
if c.SetParams(conf) {
if len(c.pid) > 0 {
c.PushPid()
return true
}
}
return false
}
func (c *CgroupTask) GroupName() string {
if len(c.name) > 0 {
return c.name
}
return c.pid
}
func (c *CgroupTask) TaskPath() string {
return path.Join(manager.GetCgroupPath(), c.GroupName())
}
func (c *CgroupTask) TaskIsExist() bool {
return fileutils2.Exists(c.TaskPath())
}
func (c *CgroupTask) createTask() bool {
if err := os.Mkdir(c.TaskPath(), os.ModePerm); err != nil {
log.Errorf("cgroup path create failed %s", err)
return false
}
return true
}
func (c *CgroupTask) ensureTask() bool {
if !c.TaskIsExist() {
if !c.createTask() {
return false
}
}
if len(c.threadIds) > 0 {
// switch to threaded mode
if !c.SetParam(CGROUP_TYPE, CGROUP_TYPE_THREADED) {
return false
}
}
if !c.SetParam(CGROUP_SUBTREE_CONTROL, fmt.Sprintf("+%s", c.hand.Module())) {
return false
}
return true
}
func (c *CgroupTask) SetParam(name, value string) bool {
err := setParam(name, value, c.GroupName())
if err != nil {
log.Errorf("Fail to set %s=%s for %s: %s", name, value, c.GroupName(), err)
return false
}
return true
}
func (c *CgroupTask) SetParams(conf map[string]string) bool {
for k, v := range conf {
if !c.SetParam(k, v) {
return false
}
}
return true
}
func (c CgroupTask) PushPid() {
if c.pid == "" {
return
}
if len(c.threadIds) > 0 {
for i := range c.threadIds {
subdir := fmt.Sprintf("/proc/%s/task/%s", c.pid, c.threadIds[i])
if fi, err := os.Stat(subdir); err != nil {
log.Errorf("Fail to stat %s in task %s", subdir, err)
continue
} else if fi.Mode().IsDir() {
stat, err := fileutils2.FileGetContents(path.Join(subdir, "stat"))
if err != nil {
log.Errorf("Fail to stat %s in task %s", stat, err)
continue
}
re := regexp.MustCompile(`\s+`)
data := re.Split(stat, -1)
if data[2] != "Z" {
c.SetParam(CGROUP_THREADS, c.threadIds[i])
}
}
}
} else {
c.SetParam(CGROUP_PROCS, c.pid)
}
}
// cgroup cpu.weight
// Convert cgroup v1 cpu.shares value to cgroup v2 cpu.weight
// https://github.com/kubernetes/enhancements/tree/master/keps/sig-node/2254-cgroup-v2#phase-1-convert-from-cgroups-v1-settings-to-v2
func CpuSharesToCpuWeight(cpuShares uint64) uint64 {
return uint64((((cpuShares - 2) * 9999) / 262142) + 1)
}
type CGroupCPUTask struct {
*CgroupTask
weight uint64
}
func (c *CGroupCPUTask) Module() string {
return "cpu"
}
func (c *CGroupCPUTask) GetConfig() map[string]string {
return map[string]string{CPU_WEIGHT: fmt.Sprintf("%d", c.weight)}
}
func (m *cgroupManager) NewCGroupCPUTask(pid, name string, cpuShares int) cgroup.ICGroupTask {
task := &CGroupCPUTask{
CgroupTask: NewCGroupBaseTask(pid, name, nil),
weight: CpuSharesToCpuWeight(uint64(cpuShares)),
}
task.SetHand(task)
return task
}
// cgroup cpuset.cpus
const (
CPUSET_CPUS = "cpuset.cpus"
CPUSET_MEMS = "cpuset.mems"
)
type CGroupCPUSetTask struct {
*CgroupTask
cpuset string
mems string
}
func (c *CGroupCPUSetTask) Module() string {
return "cpuset"
}
func (c *CGroupCPUSetTask) GetConfig() map[string]string {
if c.cpuset == "" {
c.cpuset = getParam(CPUSET_CPUS, c.GetParentGroup())
}
if c.mems == "" {
c.mems = getParam(CPUSET_MEMS, c.GetParentGroup())
}
config := map[string]string{}
if c.cpuset != "" {
config[CPUSET_CPUS] = c.cpuset
}
if c.mems != "" {
config[CPUSET_MEMS] = c.mems
}
return config
}
func (m *cgroupManager) NewCGroupCPUSetTask(pid, name, cpuset, mems string) cgroup.ICGroupTask {
task := &CGroupCPUSetTask{
CgroupTask: NewCGroupBaseTask(pid, name, nil),
cpuset: cpuset,
mems: mems,
}
task.SetHand(task)
return task
}
func (m *cgroupManager) NewCGroupSubCPUSetTask(pid, name string, cpuset string, threadIds []string) cgroup.ICGroupTask {
task := &CGroupCPUSetTask{
CgroupTask: NewCGroupBaseTask(pid, name, threadIds),
cpuset: cpuset,
}
task.SetHand(task)
return task
}
+1
View File
@@ -0,0 +1 @@
package cgroupv2 // import "yunion.io/x/onecloud/pkg/util/cgrouputils/cgroupv2"
+167
View File
@@ -0,0 +1,167 @@
// 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 cgroupv2
import (
"fmt"
"io/ioutil"
"os"
"path"
"regexp"
"strings"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/util/cgrouputils/cgroup"
"yunion.io/x/onecloud/pkg/util/fileutils2"
)
var (
manager *cgroupManager
)
type cgroupManager struct {
cgroupPath string
ioScheduler string
}
func (m *cgroupManager) GetCgroupPath() string {
return m.cgroupPath
}
func (m *cgroupManager) GetSubModulePath(module string) string {
return m.cgroupPath
}
func (m *cgroupManager) GetCgroupVersion() string {
return cgroup.CGROUP_V2
}
func (m *cgroupManager) GetIoScheduler() string {
return m.ioScheduler
}
func setParam(name, value string, groups ...string) error {
groupPath := manager.GetCgroupPath()
if len(groups) > 0 {
groups = append([]string{groupPath}, groups...)
groupPath = path.Join(groups...)
}
configPath := path.Join(groupPath, name)
return ioutil.WriteFile(configPath, []byte(value), 0644)
}
func getParam(name, group string) string {
configPath := path.Join(manager.GetCgroupPath(), group, name)
param, err := fileutils2.FileGetContents(configPath)
if err != nil {
log.Errorf("failed get cgroup config %s: %s", configPath, err)
return ""
}
return strings.TrimSpace(param)
}
func Init(cgroupPath, ioScheduler string) (*cgroupManager, error) {
if manager != nil {
return manager, nil
}
manager = &cgroupManager{
cgroupPath: cgroupPath,
ioScheduler: ioScheduler,
}
for _, module := range []string{"cpu", "cpuset"} {
err := initSubGroups(module)
if err != nil {
manager = nil
return nil, err
}
}
return manager, nil
}
func initSubGroups(module string, groups ...string) error {
err := setParam(CGROUP_SUBTREE_CONTROL, fmt.Sprintf("+%s", module), groups...)
if err != nil {
return errors.Wrapf(err, "failed add %s to cgroup.subtree_control", module)
}
return nil
}
func (m *cgroupManager) RebalanceProcesses(pids []string) {}
func (m *cgroupManager) CgroupDestroy(pid, name string) bool {
tasks := []cgroup.ICGroupTask{
&CGroupCPUTask{CgroupTask: NewCGroupBaseTask(pid, name, nil)},
&CGroupCPUSetTask{CgroupTask: NewCGroupBaseTask(pid, name, nil)},
}
for _, t := range tasks {
t.SetHand(t)
if !t.RemoveTask() {
return false
}
}
return true
}
func (m *cgroupManager) CgroupCleanAll(subName string) {
cgroupPath := manager.GetCgroupPath()
if subName != "" {
cgroupPath = path.Join(cgroupPath, subName)
}
files, err := ioutil.ReadDir(cgroupPath)
if err != nil {
log.Errorf("%s GetTasks failed: %s", cgroupPath, err)
return
}
ids := []string{}
for _, file := range files {
ids = append(ids, file.Name())
}
re1 := regexp.MustCompile(`^\d+$`)
re2 := regexp.MustCompile(`^server_[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}_\d+$`)
for _, pid := range ids {
spid := ""
if re1.MatchString(pid) {
spid = pid
} else if re2.MatchString(pid) {
segs := strings.Split(pid, "_")
spid = segs[len(segs)-1]
}
if fileutils2.IsDir(path.Join(cgroupPath, pid)) {
if !fileutils2.Exists(path.Join("/proc", spid)) {
log.Infof("Cgroup cleanup %s", path.Join(cgroupPath, pid))
}
subFiles, err := ioutil.ReadDir(path.Join(cgroupPath, pid))
if err != nil {
log.Errorf("sub dir %s GetTaskIds failed: %s", path.Join(cgroupPath, pid), err)
} else {
for _, fi := range subFiles {
if !fi.IsDir() {
continue
}
if err := os.Remove(path.Join(cgroupPath, pid, fi.Name())); err != nil {
log.Errorf("CgroupCleanAll pid=%s tid=%s error: %s", pid, fi.Name(), err)
}
}
}
if err := os.Remove(path.Join(cgroupPath, pid)); err != nil {
log.Errorf("CgroupCleanAll pid=%s error: %s", pid, err)
}
}
}
}
+117
View File
@@ -0,0 +1,117 @@
// 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 cgrouputils
import (
"strings"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/util/cgrouputils/cgroup"
"yunion.io/x/onecloud/pkg/util/cgrouputils/cgroupv1"
"yunion.io/x/onecloud/pkg/util/cgrouputils/cgroupv2"
"yunion.io/x/onecloud/pkg/util/fileutils2"
"yunion.io/x/onecloud/pkg/util/procutils"
)
const (
CGROUP_PATH_SYSFS = "/sys/fs/cgroup"
CGROUP_PATH_ROOT = "/cgroup"
CPUSET_SCHED_LOAD_BALANCE = "cpuset.sched_load_balance"
CPUSET_CLONE_CHILDREN = "cgroup.clone_children"
)
type ICgroupManager interface {
RebalanceProcesses(pids []string)
CgroupCleanAll(subName string)
CgroupDestroy(pid, name string) bool
GetCgroupPath() string
GetSubModulePath(module string) string
NewCGroupCPUSetTask(pid, name, cpuset, mems string) cgroup.ICGroupTask
NewCGroupCPUTask(pid, name string, cpuShares int) cgroup.ICGroupTask
NewCGroupSubCPUSetTask(pid, name string, cpuset string, threadIds []string) cgroup.ICGroupTask
}
func GetCgroupVersion() string {
return cgroupManager.GetCgroupPath()
}
func RebalanceProcesses(pids []string) {
cgroupManager.RebalanceProcesses(pids)
}
func CgroupCleanAll(subName string) {
cgroupManager.CgroupCleanAll(subName)
}
func CgroupDestroy(pid, name string) bool {
return cgroupManager.CgroupDestroy(pid, name)
}
func GetCgroupPath() string {
return cgroupManager.GetCgroupPath()
}
func GetSubModulePath(module string) string {
return cgroupManager.GetSubModulePath(module)
}
func NewCGroupCPUTask(pid, name string, cpuShares int) cgroup.ICGroupTask {
return cgroupManager.NewCGroupCPUTask(pid, name, cpuShares)
}
func NewCGroupCPUSetTask(pid, name, cpuset, mems string) cgroup.ICGroupTask {
return cgroupManager.NewCGroupCPUSetTask(pid, name, cpuset, mems)
}
func NewCGroupSubCPUSetTask(pid, name string, cpuset string, threadIds []string) cgroup.ICGroupTask {
return cgroupManager.NewCGroupSubCPUSetTask(pid, name, cpuset, threadIds)
}
var cgroupManager ICgroupManager
func Init(ioScheduler string) error {
if cgroupManager != nil {
return nil
}
cgroupPath := ""
if fileutils2.Exists(CGROUP_PATH_SYSFS) {
cgroupPath = CGROUP_PATH_SYSFS
} else if fileutils2.Exists("CGROUP_PATH_ROOT") {
cgroupPath = CGROUP_PATH_ROOT
}
if cgroupPath == "" {
return errors.Errorf("Can't detect cgroup path")
}
output, err := procutils.NewCommand("stat", "-fc", "%T", cgroupPath).Output()
if err != nil {
return errors.Wrapf(err, "stat cgroup path %s", cgroupPath)
}
cgroupfs := strings.TrimSpace(string(output))
if cgroupfs == "cgroup2fs" {
// cgroup v2
cgroupManager, err = cgroupv2.Init(cgroupPath, ioScheduler)
} else {
// cgroup v1
cgroupManager, err = cgroupv1.Init(cgroupPath, ioScheduler)
}
if err != nil {
return errors.Wrap(err, "init cgroup")
}
return nil
}