Merge pull request #1245 in YUNIONIO/onecloud from ~WANYAOQI/onecloud:feature/wyq/host-server-v2 to release/2.8.0

* commit 'b67eb79797040b211d0b2ab4bb2f8988d06c78d4':
  fix cpu set
  fix cgroups
This commit is contained in:
邱剑
2019-03-18 22:08:27 +08:00
7 changed files with 148 additions and 103 deletions
+9 -1
View File
@@ -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 {
+9 -1
View File
@@ -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)
}
}()
}
+5 -28
View File
@@ -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 {
+1 -1
View File
@@ -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, "", "")
}
+26 -24
View File
@@ -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 ""
}
@@ -424,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
}
/**
@@ -449,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
}
/**
@@ -480,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
}
/**
@@ -514,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 {
@@ -536,9 +541,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 +566,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
}
+68 -48
View File
@@ -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
}
+30
View File
@@ -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)
}