From 44f9b4c818e8e24449ad732cb6d65fff732009b7 Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Wed, 13 Mar 2019 13:54:11 +0800 Subject: [PATCH 1/2] fix cgroups --- pkg/hostman/guestman/qemu-kvm.go | 33 +++++------------------------ pkg/util/cgrouputils/cgrouputils.go | 30 +++++++++++--------------- 2 files changed, 17 insertions(+), 46 deletions(-) diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index c2e9364916..b30030d91b 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -995,35 +995,12 @@ func (s *SKVMGuestInstance) setCgroupIo() { } func (s *SKVMGuestInstance) setCgroupCpu() { - cpu, _ := s.Desc.Int("cpu") - cgrouputils.CgroupSet(strconv.Itoa(s.cgroupPid), int(cpu)) + var ( + cpu, _ = s.Desc.Int("cpu") + cpuWeight = 1024 + ) - // TODO XXX - /* - var ( - cpuWeight = 1024 - cpuPeriod = 0 - cpuQuota = 0 - appTags = s.getApptags() - meta, _ = s.Desc.Get("metadata") - ) - - if meta != nil { - if meta.Contains("__cpu_weight") { - cpuWeight, _ = meta.Int("__cpu_weight") - } - if meta.Contains("__cpu_period") { - cpuPeriod, _ = meta.Int("__cpu_period") - } else { - cpuPeriod = -1 - } - if meta.Contains("__cpu_quota") { - cpuQuota, _ = meta.Int("__cpu_quota") - } else { - cpuQuota = -1 - } - } - */ + cgrouputils.CgroupSet(strconv.Itoa(s.cgroupPid), int(cpu)*cpuWeight) } func (s *SKVMGuestInstance) CreateFromDesc(desc jsonutils.JSONObject) error { diff --git a/pkg/util/cgrouputils/cgrouputils.go b/pkg/util/cgrouputils/cgrouputils.go index 1a71767992..79fbd9ae5c 100644 --- a/pkg/util/cgrouputils/cgrouputils.go +++ b/pkg/util/cgrouputils/cgrouputils.go @@ -30,6 +30,7 @@ var ( ) type ICGroupTask interface { + InitTask(hand ICGroupTask, coreNum int, pid string) SetPid(string) SetWeight(coreNum int) SetHand(hand ICGroupTask) @@ -108,19 +109,9 @@ func GetRootParam(module, name, pid string) string { } func SetRootParam(module, name, value, pid string) bool { - if param := GetRootParam(module, name, pid); param != value { - fi, err := os.Open(GetTaskParamPath(module, name, pid)) - if err == nil { - _, err = fi.Write([]byte(value)) - if err != nil { - err = fi.Close() - } else { - log.Errorln(err) - } - } else { - log.Errorln(err) - } - + 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" @@ -168,6 +159,12 @@ func (c *CGroupTask) SetPid(pid string) { c.pid = pid } +func (c *CGroupTask) InitTask(hand ICGroupTask, coreNum int, pid string) { + c.SetHand(hand) + c.SetWeight(coreNum) + c.SetPid(pid) +} + func (c *CGroupTask) Module() string { return "" } @@ -536,9 +533,7 @@ func CgroupSet(pid string, coreNum int) bool { &CGroupMemoryTask{&CGroupTask{}}, } for _, hand := range tasks { - hand.SetHand(hand) - hand.SetPid(pid) - hand.SetWeight(coreNum) + hand.InitTask(hand, coreNum, pid) if !hand.SetTask() { return false } @@ -563,8 +558,7 @@ func CgroupDestroy(pid string) bool { &CGroupIOHardlimitTask{CGroupIOTask: &CGroupIOTask{&CGroupTask{}}}, } for _, hand := range tasks { - hand.SetHand(hand) - hand.SetPid(pid) + hand.InitTask(hand, 0, pid) if !hand.RemoveTask() { return false } From b67eb79797040b211d0b2ab4bb2f8988d06c78d4 Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Wed, 13 Mar 2019 21:07:03 +0800 Subject: [PATCH 2/2] fix cpu set --- pkg/appsrv/appsrv.go | 10 +- pkg/hostman/guestman/guestman.go | 10 +- pkg/hostman/metadata/metadatahandler.go | 2 +- pkg/util/cgrouputils/cgrouputils.go | 20 ++-- pkg/util/cgrouputils/cpuutils.go | 116 ++++++++++++++---------- pkg/util/cgrouputils/cpuutils_test.go | 30 ++++++ 6 files changed, 131 insertions(+), 57 deletions(-) create mode 100644 pkg/util/cgrouputils/cpuutils_test.go diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index ef8cfbf227..4fdae4346b 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -397,7 +397,16 @@ func (app *Application) ListenAndServeWithCleanup(addr string, onStop func()) { func (app *Application) ListenAndServeTLSWithCleanup(addr string, certFile, keyFile string, onStop func()) { s := app.initServer(addr) app.registerCleanShutdown(s, onStop) + app.listenAndServe(s, certFile, keyFile) + app.waitCleanShutdown() +} +func (app *Application) ListenAndServeWithoutCleanup(addr, certFile, keyFile string) { + s := app.initServer(addr) + app.listenAndServe(s, certFile, keyFile) +} + +func (app *Application) listenAndServe(s *http.Server, certFile, keyFile string) { var err error if len(certFile) == 0 && len(keyFile) == 0 { err = s.ListenAndServe() @@ -407,7 +416,6 @@ func (app *Application) ListenAndServeTLSWithCleanup(addr string, certFile, keyF if err != nil && err != http.ErrServerClosed { log.Fatalf("ListAndServer fail: %s", err) } - app.waitCleanShutdown() } func isJsonContentType(r *http.Request) bool { diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 344c60b03d..9b45c9ebde 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -6,6 +6,7 @@ import ( "io/ioutil" "os" "path" + "runtime/debug" "strings" "sync" "time" @@ -156,11 +157,18 @@ func (m *SGuestManager) StartCpusetBalancer() { return } go func() { + defer func() { + if r := recover(); r != nil { + debug.PrintStack() + log.Errorf("Cpuset balancer failed %s", r) + } + }() for { + time.Sleep(time.Second * 120) + if options.HostOptions.EnableCpuBinding { m.cpusetBalance() } - time.Sleep(time.Second * 120) } }() } diff --git a/pkg/hostman/metadata/metadatahandler.go b/pkg/hostman/metadata/metadatahandler.go index 520894047f..82ef69935c 100644 --- a/pkg/hostman/metadata/metadatahandler.go +++ b/pkg/hostman/metadata/metadatahandler.go @@ -251,5 +251,5 @@ func StartService(app *appsrv.Application, address string, port int) { addMetadataHandler("", app) addr := net.JoinHostPort(address, strconv.Itoa(port)) log.Infof("Host Metadata Start listen on %s://%s", "http", addr) - app.ListenAndServeWithCleanup(addr, nil) + app.ListenAndServeWithoutCleanup(addr, "", "") } diff --git a/pkg/util/cgrouputils/cgrouputils.go b/pkg/util/cgrouputils/cgrouputils.go index 79fbd9ae5c..1cf340d132 100644 --- a/pkg/util/cgrouputils/cgrouputils.go +++ b/pkg/util/cgrouputils/cgrouputils.go @@ -421,7 +421,9 @@ func (c *CGroupIOTask) init() bool { } func NewCGroupIOTask(pid string, coreNum int) *CGroupIOTask { - return &CGroupIOTask{NewCGroupTask(pid, coreNum)} + task := &CGroupIOTask{NewCGroupTask(pid, coreNum)} + task.SetHand(task) + return task } /** @@ -446,13 +448,15 @@ func (c *CGroupIOHardlimitTask) GetConfig() map[string]string { return config } -func NewCGroupIOHardlimitTask(pid string, mem int, params map[string]int, devId string) CGroupIOHardlimitTask { - return CGroupIOHardlimitTask{ +func NewCGroupIOHardlimitTask(pid string, mem int, params map[string]int, devId string) *CGroupIOHardlimitTask { + task := &CGroupIOHardlimitTask{ CGroupIOTask: NewCGroupIOTask(pid, 0), cpuNum: mem, params: params, devId: devId, } + task.SetHand(task) + return task } /** @@ -477,10 +481,12 @@ func (c *CGroupMemoryTask) GetConfig() map[string]string { return map[string]string{MEMORY_SWAPPINESS: fmt.Sprintf("%d", vm_swappiness)} } -func NewCGroupMemoryTask(pid string, coreNum int) CGroupMemoryTask { - return CGroupMemoryTask{ +func NewCGroupMemoryTask(pid string, coreNum int) *CGroupMemoryTask { + task := &CGroupMemoryTask{ CGroupTask: NewCGroupTask(pid, coreNum), } + task.SetHand(task) + return task } /** @@ -511,10 +517,12 @@ func (c *CGroupCPUSetTask) GetConfig() map[string]string { } func NewCGroupCPUSetTask(pid string, coreNum int, cpuset string) CGroupCPUSetTask { - return CGroupCPUSetTask{ + task := CGroupCPUSetTask{ CGroupTask: NewCGroupTask(pid, coreNum), cpuset: cpuset, } + task.SetHand(&task) + return task } func Init() bool { diff --git a/pkg/util/cgrouputils/cpuutils.go b/pkg/util/cgrouputils/cpuutils.go index 523a9e80d6..e3d820ed3d 100644 --- a/pkg/util/cgrouputils/cpuutils.go +++ b/pkg/util/cgrouputils/cpuutils.go @@ -21,6 +21,8 @@ type CPU struct { func NewCPU() (*CPU, error) { var cpu = new(CPU) + cpu.DieList = make([]*CPUDie, 0) + cpuinfo, err := fileutils2.FileGetContents("/proc/cpuinfo") if err != nil { return nil, err @@ -30,23 +32,30 @@ func NewCPU() (*CPU, error) { for _, line := range strings.Split(cpuinfo, "\n") { parts := strings.Split(line, ":") if len(parts) == 2 { - if parts[0] == "processor" { - val, err := strconv.Atoi(parts[1]) + key := strings.TrimSpace(parts[0]) + val := strings.TrimSpace(parts[1]) + if key == "processor" { + iVal, err := strconv.Atoi(val) if err != nil { return nil, err } - core = NewCPUCore(val) - cpu.AddCore(core) + if core != nil { + cpu.AddCore(core) + } + core = NewCPUCore(iVal) } else if core != nil { - core.setInfo(parts[0], parts[1]) + core.setInfo(key, val) } } } + if core != nil { + cpu.AddCore(core) + } return cpu, nil } func (c *CPU) AddCore(core *CPUCore) { - for len(c.DieList) < core.PhysicalId { + for len(c.DieList) < core.PhysicalId+1 { c.DieList = append(c.DieList, NewCPUDie(len(c.DieList))) } c.DieList[core.PhysicalId].AddCore(core) @@ -176,8 +185,11 @@ func Average(arr []float64) float64 { return total / float64(len(arr)) } -func GetProcessWeight(share int, util float64) float64 { - return float64(share) * (util*0.8 + 30) +func GetProcessWeight(share *int, util float64) float64 { + if share != nil { + return float64(*share) * (util*0.8 + 30) + } + return 0.0 } func NewProcessCPUinfo(pid int) (*ProcessCPUinfo, error) { @@ -185,52 +197,60 @@ func NewProcessCPUinfo(pid int) (*ProcessCPUinfo, error) { cpuinfo.Pid = pid spid := strconv.Itoa(pid) - share := NewCGroupCPUTask(spid, 0).GetParam("cpu.shares") - ishare, err := strconv.Atoi(share) - if err != nil { - log.Errorln(err) - } else { - if ishare != 0 { - cpuinfo.Share = &ishare - proc, err := process.NewProcess(int32(pid)) - if err != nil { - log.Errorln(err) - return nil, err - } - util, err := proc.CPUPercent() - if err != nil { - log.Errorln(err) - return nil, err - } - util /= float64(ishare) - uHistory := FetchHistoryUtil() - - var utils = []float64{} - if _, ok := uHistory[spid]; ok { - utils = uHistory[spid] - } - utils = append(utils, util) - for len(utils) > MAX_HISTORY_UTIL_COUNT { - utils = utils[1:] - } - uHistory[spid] = utils - - cpuinfo.Util = Average(utils) - } - } - - cpuset := NewCGroupCPUSetTask(fmt.Sprintf("%s", pid), 0, "").GetParam("cpuset.cpus") - if len(cpuset) > 0 { - c, err := GetSystemCpu() + cpuTask := NewCGroupCPUTask(spid, 0) + if cpuTask.taskIsExist() { + share := cpuTask.GetParam("cpu.shares") + ishare, err := strconv.Atoi(share) + ishare /= 1024.0 if err != nil { log.Errorln(err) } else { - icpuset := c.GetPhysicalId(ParseCpusetStr(cpuset)) - cpuinfo.Cpuset = &icpuset + cpuinfo.Share = &ishare } } - cpuinfo.Weight = GetProcessWeight(*cpuinfo.Share, cpuinfo.Util) + cpusetTask := NewCGroupCPUSetTask(fmt.Sprintf("%d", pid), 0, "") + if cpusetTask.taskIsExist() { + cpuset := cpusetTask.GetParam("cpuset.cpus") + if len(cpuset) > 0 { + c, err := GetSystemCpu() + if err != nil { + log.Errorln(err) + } else { + icpuset := c.GetPhysicalId(ParseCpusetStr(cpuset)) + cpuinfo.Cpuset = &icpuset + } + } + } + + if cpuinfo.Share != nil { + proc, err := process.NewProcess(int32(pid)) + if err != nil { + log.Errorln(err) + return nil, err + } + util, err := proc.CPUPercent() + if err != nil { + log.Errorln(err) + return nil, err + } + util /= float64(*cpuinfo.Share) + uHistory := FetchHistoryUtil() + + var utils = []float64{} + if _, ok := uHistory[spid]; ok { + utils = uHistory[spid] + } + utils = append(utils, util) + for len(utils) > MAX_HISTORY_UTIL_COUNT { + utils = utils[1:] + } + uHistory[spid] = utils + + cpuinfo.Util = Average(utils) + } + + cpuinfo.Weight = GetProcessWeight(cpuinfo.Share, cpuinfo.Util) return cpuinfo, nil } diff --git a/pkg/util/cgrouputils/cpuutils_test.go b/pkg/util/cgrouputils/cpuutils_test.go new file mode 100644 index 0000000000..f753bfa1cd --- /dev/null +++ b/pkg/util/cgrouputils/cpuutils_test.go @@ -0,0 +1,30 @@ +package cgrouputils + +import ( + "os" + "testing" +) + +func TestCPUUtils(t *testing.T) { + cpu, err := GetSystemCpu() + if err != nil { + t.Errorf("NewCPU error %s", err) + } + + t.Logf("die list %v", cpu.DieList) + t.Logf("physical cpu num %d", cpu.GetPhysicalNum()) + + cpusetStr := cpu.GetCpuset(0) + t.Logf("cpu 0 set str %s", cpusetStr) + physicalNum := cpu.GetPhysicalId(cpusetStr) + if physicalNum != 0 { + t.Errorf("failed get physical id %d from cpuset %s", 0, cpusetStr) + } + + pid := os.Getpid() + proc, err := NewProcessCPUinfo(pid) + if err != nil { + t.Errorf("New process cpu info error %s", err) + } + t.Logf("Shared %v %v %f %f", proc.Share, proc.Cpuset, proc.Util, proc.Weight) +}