diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index c35fd156e5..1d819b5a2c 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -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") } diff --git a/pkg/hostman/hostinfo/hostinfo.go b/pkg/hostman/hostinfo/hostinfo.go index 1419b555d3..bb1c637f46 100644 --- a/pkg/hostman/hostinfo/hostinfo.go +++ b/pkg/hostman/hostinfo/hostinfo.go @@ -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") { diff --git a/pkg/hostman/hostinfo/hostinfohelper.go b/pkg/hostman/hostinfo/hostinfohelper.go index b62a5b1c5a..705162e411 100644 --- a/pkg/hostman/hostinfo/hostinfohelper.go +++ b/pkg/hostman/hostinfo/hostinfohelper.go @@ -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"` diff --git a/pkg/mcclient/options/compute/servers.go b/pkg/mcclient/options/compute/servers.go index ef98dcb653..fcac143768 100644 --- a/pkg/mcclient/options/compute/servers.go +++ b/pkg/mcclient/options/compute/servers.go @@ -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 } diff --git a/pkg/util/cgrouputils/cgroup/doc.go b/pkg/util/cgrouputils/cgroup/doc.go new file mode 100644 index 0000000000..db03c717c0 --- /dev/null +++ b/pkg/util/cgrouputils/cgroup/doc.go @@ -0,0 +1 @@ +package cgroup // import "yunion.io/x/onecloud/pkg/util/cgrouputils/cgroup" diff --git a/pkg/util/cgrouputils/cgroup/interface.go b/pkg/util/cgrouputils/cgroup/interface.go new file mode 100644 index 0000000000..866721083c --- /dev/null +++ b/pkg/util/cgrouputils/cgroup/interface.go @@ -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 +} diff --git a/pkg/util/cgrouputils/cgrouputils.go b/pkg/util/cgrouputils/cgroupv1/cgrouputils.go similarity index 53% rename from pkg/util/cgrouputils/cgrouputils.go rename to pkg/util/cgrouputils/cgroupv1/cgrouputils.go index bd132fed3b..4daa8622d3 100644 --- a/pkg/util/cgrouputils/cgrouputils.go +++ b/pkg/util/cgrouputils/cgroupv1/cgrouputils.go @@ -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) - } -} diff --git a/pkg/util/cgrouputils/cgrouputils_test.go b/pkg/util/cgrouputils/cgroupv1/cgrouputils_test.go similarity index 87% rename from pkg/util/cgrouputils/cgrouputils_test.go rename to pkg/util/cgrouputils/cgroupv1/cgrouputils_test.go index 5324c29418..284997c218 100644 --- a/pkg/util/cgrouputils/cgrouputils_test.go +++ b/pkg/util/cgrouputils/cgroupv1/cgrouputils_test.go @@ -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("") } diff --git a/pkg/util/cgrouputils/cpusetutils.go b/pkg/util/cgrouputils/cgroupv1/cpusetutils.go similarity index 97% rename from pkg/util/cgrouputils/cpusetutils.go rename to pkg/util/cgrouputils/cgroupv1/cpusetutils.go index 884a93d910..b87d6571b3 100644 --- a/pkg/util/cgrouputils/cpusetutils.go +++ b/pkg/util/cgrouputils/cgroupv1/cpusetutils.go @@ -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() } } diff --git a/pkg/util/cgrouputils/cpuutils.go b/pkg/util/cgrouputils/cgroupv1/cpuutils.go similarity index 96% rename from pkg/util/cgrouputils/cpuutils.go rename to pkg/util/cgrouputils/cgroupv1/cpuutils.go index 94b304e961..e252bf5a64 100644 --- a/pkg/util/cgrouputils/cpuutils.go +++ b/pkg/util/cgrouputils/cgroupv1/cpuutils.go @@ -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() diff --git a/pkg/util/cgrouputils/cpuutils_test.go b/pkg/util/cgrouputils/cgroupv1/cpuutils_test.go similarity index 98% rename from pkg/util/cgrouputils/cpuutils_test.go rename to pkg/util/cgrouputils/cgroupv1/cpuutils_test.go index 37875455c7..7de05e222a 100644 --- a/pkg/util/cgrouputils/cpuutils_test.go +++ b/pkg/util/cgrouputils/cgroupv1/cpuutils_test.go @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -package cgrouputils +package cgroupv1 import ( "os" diff --git a/pkg/util/cgrouputils/cgroupv1/doc.go b/pkg/util/cgrouputils/cgroupv1/doc.go new file mode 100644 index 0000000000..44691d57fd --- /dev/null +++ b/pkg/util/cgrouputils/cgroupv1/doc.go @@ -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" diff --git a/pkg/util/cgrouputils/cgroupv1/manager.go b/pkg/util/cgrouputils/cgroupv1/manager.go new file mode 100644 index 0000000000..dfd7556944 --- /dev/null +++ b/pkg/util/cgrouputils/cgroupv1/manager.go @@ -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) + } + } + } + } +} diff --git a/pkg/util/cgrouputils/cgroupv2/cgroup.go b/pkg/util/cgrouputils/cgroupv2/cgroup.go new file mode 100644 index 0000000000..740370ecc4 --- /dev/null +++ b/pkg/util/cgrouputils/cgroupv2/cgroup.go @@ -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 +} diff --git a/pkg/util/cgrouputils/cgroupv2/doc.go b/pkg/util/cgrouputils/cgroupv2/doc.go new file mode 100644 index 0000000000..b7133d67a4 --- /dev/null +++ b/pkg/util/cgrouputils/cgroupv2/doc.go @@ -0,0 +1 @@ +package cgroupv2 // import "yunion.io/x/onecloud/pkg/util/cgrouputils/cgroupv2" diff --git a/pkg/util/cgrouputils/cgroupv2/manager.go b/pkg/util/cgrouputils/cgroupv2/manager.go new file mode 100644 index 0000000000..5236cb4f2b --- /dev/null +++ b/pkg/util/cgrouputils/cgroupv2/manager.go @@ -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) + } + } + } +} diff --git a/pkg/util/cgrouputils/manager.go b/pkg/util/cgrouputils/manager.go new file mode 100644 index 0000000000..2907098332 --- /dev/null +++ b/pkg/util/cgrouputils/manager.go @@ -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 +}