From 46e85d5249f7f3021895dde93de65be5c01b970f Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Wed, 27 Feb 2019 18:35:01 +0800 Subject: [PATCH] kvm guest support hotplug cpu mem, fix hmp monitor --- pkg/cloudcommon/db/taskman/tasks.go | 10 +- pkg/compute/guestdrivers/base.go | 4 + pkg/compute/guestdrivers/kvm.go | 56 +++++++- pkg/compute/models/guest_actions.go | 7 + pkg/compute/models/guestdrivers.go | 2 + pkg/compute/tasks/guest_change_config_task.go | 19 +-- pkg/hostman/guesthandlers/guesthandler.go | 19 ++- pkg/hostman/guestman/guesthelper.go | 6 + pkg/hostman/guestman/guestman.go | 15 ++ pkg/hostman/guestman/guesttasks.go | 132 ++++++++++++++++++ pkg/hostman/guestman/qemu-kvmhelper.go | 4 +- pkg/hostman/hostutils/hostutils.go | 8 ++ pkg/hostman/monitor/hmp.go | 71 ++++++++-- pkg/hostman/monitor/hmp_test.go | 17 +-- pkg/hostman/monitor/monitor.go | 6 + pkg/hostman/monitor/qmp.go | 57 ++++++++ pkg/mcclient/modules/mod_tasks.go | 8 +- 17 files changed, 399 insertions(+), 42 deletions(-) diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index e9c703ac6a..e5bc6a2d0b 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -317,7 +317,6 @@ func (manager *STaskManager) execTask(taskId string, data jsonutils.JSONObject) } func execITask(taskValue reflect.Value, task *STask, odata jsonutils.JSONObject, isMulti bool) { - var err error ctxData := task.GetRequestContext() ctx := ctxData.GetContext() @@ -328,9 +327,12 @@ func execITask(taskValue reflect.Value, task *STask, odata jsonutils.JSONObject, taskStatus, _ := data.GetString("__status__") if len(taskStatus) > 0 && taskStatus != "OK" { taskFailed = true - data, err = data.Get("__reason__") - if err != nil { - data = jsonutils.NewString(fmt.Sprintf("Task failed due to unknown remote errors! %s", odata)) + if vdata, ok := data.(*jsonutils.JSONDict); ok { + reason, err := vdata.Get("__reason__") // only dict support Get + if err != nil { + reason = jsonutils.NewString(fmt.Sprintf("Task failed due to unknown remote errors! %s", odata)) + vdata.Set("__reason__", reason) + } } } } else { diff --git a/pkg/compute/guestdrivers/base.go b/pkg/compute/guestdrivers/base.go index 5321e6c8c2..dd796c1506 100644 --- a/pkg/compute/guestdrivers/base.go +++ b/pkg/compute/guestdrivers/base.go @@ -235,3 +235,7 @@ func (self *SBaseGuestDriver) IsSupportEip() bool { func (self *SBaseGuestDriver) NeedStopForChangeSpec() bool { return true } + +func (self *SBaseGuestDriver) OnGuestChangeCpuMemFailed(ctx context.Context, guest *models.SGuest, data *jsonutils.JSONDict, task taskman.ITask) error { + return nil +} diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index 024f563727..de33754023 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -11,11 +11,13 @@ import ( "yunion.io/x/log" "yunion.io/x/pkg/utils" + "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/compute/options" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/util/httputils" + "yunion.io/x/onecloud/pkg/util/logclient" ) type SKVMGuestDriver struct { @@ -255,10 +257,28 @@ func (self *SKVMGuestDriver) OnDeleteGuestFinalCleanup(ctx context.Context, gues return nil } +func (self *SKVMGuestDriver) NeedStopForChangeSpec() bool { + return false +} + func (self *SKVMGuestDriver) RequestChangeVmConfig(ctx context.Context, guest *models.SGuest, task taskman.ITask, instanceType string, vcpuCount, vmemSize int64) error { - // pass - task.ScheduleRun(nil) - return nil + if jsonutils.QueryBoolean(task.GetParams(), "guest_online", false) { + header := task.GetTaskRequestHeader() + body := jsonutils.NewDict() + if vcpuCount > int64(guest.VcpuCount) { + body.Set("add_cpu", jsonutils.NewInt(vcpuCount-int64(guest.VcpuCount))) + } + if vmemSize > int64(guest.VmemSize) { + body.Set("add_mem", jsonutils.NewInt(vmemSize-int64(guest.VmemSize))) + } + host := guest.GetHost() + url := fmt.Sprintf("%s/servers/%s/hotplug-cpu-mem", host.ManagerUri, guest.Id) + _, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, body, false) + return err + } else { + task.ScheduleRun(nil) + return nil + } } func (self *SKVMGuestDriver) RequestSoftReset(ctx context.Context, guest *models.SGuest, task taskman.ITask) error { @@ -353,3 +373,33 @@ func (self *SKVMGuestDriver) RequestSyncToBackup(ctx context.Context, guest *mod } return nil } + +// kvm guest must add cpu first +// if body has add_cpu_failed indicate dosen't exec add mem +// 1. cpu added part of request --> add_cpu_failed: true && added_cpu: count +// 2. cpu added all of request add mem failed --> add_mem_failed: true +func (self *SKVMGuestDriver) OnGuestChangeCpuMemFailed(ctx context.Context, guest *models.SGuest, data *jsonutils.JSONDict, task taskman.ITask) error { + var cpuAdded int64 + if jsonutils.QueryBoolean(data, "add_cpu_failed", false) { + cpuAdded, _ = data.Int("added_cpu") + } else if jsonutils.QueryBoolean(data, "add_mem_failed", false) { + vcpuCount, _ := task.GetParams().Int("vcpu_count") + if vcpuCount-int64(guest.VcpuCount) > 0 { + cpuAdded = vcpuCount - int64(guest.VcpuCount) + } + } + if cpuAdded > 0 { + _, err := guest.GetModelManager().TableSpec().Update(guest, func() error { + guest.VcpuCount = guest.VcpuCount + int8(cpuAdded) + return nil + }) + if err != nil { + return err + } + db.OpsLog.LogEvent(guest, db.ACT_CHANGE_FLAVOR, + fmt.Sprintf("Change config task failed but added cpu count %d", cpuAdded), task.GetUserCred()) + logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_VM_CHANGE_FLAVOR, + fmt.Sprintf("Change config task failed but added cpu count %d", cpuAdded), task.GetUserCred(), false) + } + return nil +} diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 4b97f24d4d..489a41459e 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -1582,6 +1582,10 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T return nil, httperrors.NewInvalidStatusError("Not allow to change config") } + if len(self.BackupHostId) > 0 { + return nil, httperrors.NewBadRequestError("Guest have backup not allow to change config") + } + changeStatus, err := self.GetDriver().GetChangeConfigStatus() if err != nil { return nil, httperrors.NewInputParameterError(err.Error()) @@ -1749,6 +1753,9 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T if self.Status != VM_RUNNING && jsonutils.QueryBoolean(data, "auto_start", false) { confs.Add(jsonutils.NewBool(true), "auto_start") } + if self.Status == VM_RUNNING { + confs.Set("guest_online", jsonutils.JSONTrue) + } log.Debugf("%s", confs.String()) diff --git a/pkg/compute/models/guestdrivers.go b/pkg/compute/models/guestdrivers.go index fd2be736ef..5d9e07c40a 100644 --- a/pkg/compute/models/guestdrivers.go +++ b/pkg/compute/models/guestdrivers.go @@ -126,6 +126,8 @@ type IGuestDriver interface { IsSupportEip() bool NeedStopForChangeSpec() bool + + OnGuestChangeCpuMemFailed(ctx context.Context, guest *SGuest, data *jsonutils.JSONDict, task taskman.ITask) error } var guestDrivers map[string]IGuestDriver diff --git a/pkg/compute/tasks/guest_change_config_task.go b/pkg/compute/tasks/guest_change_config_task.go index b8f947677a..8a07a3630b 100644 --- a/pkg/compute/tasks/guest_change_config_task.go +++ b/pkg/compute/tasks/guest_change_config_task.go @@ -5,6 +5,7 @@ import ( "fmt" "yunion.io/x/jsonutils" + "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" @@ -24,7 +25,7 @@ func init() { func (self *GuestChangeConfigTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { _, err := self.Params.Get("resize") if err == nil { - self.SetStage("on_disks_resize_complete", nil) + self.SetStage("OnDisksResizeComplete", nil) self.OnDisksResizeComplete(ctx, obj, data) } else { guest := obj.(*models.SGuest) @@ -100,7 +101,7 @@ func (self *GuestChangeConfigTask) DoCreateDisksTask(ctx context.Context, guest return } data := (iCreateData).(*jsonutils.JSONDict) - self.SetStage("on_create_disks_complete", nil) + self.SetStage("OnCreateDisksComplete", nil) guest.StartGuestCreateDiskTask(ctx, self.UserCred, data, self.GetTaskId()) } @@ -137,11 +138,6 @@ func (self *GuestChangeConfigTask) startGuestChangeCpuMemSpec(ctx context.Contex } } -func (self *GuestChangeConfigTask) OnGuestChangeCpuMemSpecCompleteFailed(ctx context.Context, obj db.IStandaloneModel, err jsonutils.JSONObject) { - guest := obj.(*models.SGuest) - self.markStageFailed(ctx, guest, fmt.Sprintf("guest.GetDriver().RequestChangeVmConfig fail %s", err)) -} - func (self *GuestChangeConfigTask) OnGuestChangeCpuMemSpecComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { guest := obj.(*models.SGuest) @@ -200,8 +196,15 @@ func (self *GuestChangeConfigTask) OnGuestChangeCpuMemSpecComplete(ctx context.C self.OnGuestChangeCpuMemSpecFinish(ctx, guest) } +func (self *GuestChangeConfigTask) OnGuestChangeCpuMemSpecCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + if err := guest.GetDriver().OnGuestChangeCpuMemFailed(ctx, guest, data.(*jsonutils.JSONDict), self); err != nil { + log.Errorln(err) + } + self.markStageFailed(ctx, guest, fmt.Sprintf("guest.GetDriver().RequestChangeVmConfig fail %s", data)) +} + func (self *GuestChangeConfigTask) OnGuestChangeCpuMemSpecFinish(ctx context.Context, guest *models.SGuest) { - self.SetStage("on_sync_config_complete", nil) + self.SetStage("OnSyncConfigComplete", nil) err := guest.StartSyncTask(ctx, self.UserCred, false, self.GetTaskId()) if err != nil { self.markStageFailed(ctx, guest, fmt.Sprintf("StartSyncstatus fail %s", err)) diff --git a/pkg/hostman/guesthandlers/guesthandler.go b/pkg/hostman/guesthandlers/guesthandler.go index 50760596c0..fa7bdceccf 100644 --- a/pkg/hostman/guesthandlers/guesthandler.go +++ b/pkg/hostman/guesthandlers/guesthandler.go @@ -42,7 +42,8 @@ var ( "live-migrate": guestLiveMigrate, "resume": guestResume, // "start-nbd-server": guestStartNbdServer, - "drive-mirror": guestDriveMirror, + "drive-mirror": guestDriveMirror, + "hotplug-cpu-mem": guestHotplugCpuMem, } ) @@ -320,6 +321,22 @@ func guestDriveMirror(ctx context.Context, sid string, body jsonutils.JSONObject return nil, nil } +func guestHotplugCpuMem(ctx context.Context, sid string, body jsonutils.JSONObject) (interface{}, error) { + if !guestman.GetGuestManager().IsGuestExist(sid) { + return nil, httperrors.NewNotFoundError("Guest %s not found", sid) + } + + if guestman.GetGuestManager().Status(sid) != "running" { + return nil, httperrors.NewBadRequestError("Guest %s not running", sid) + } + + addCpuCount, _ := body.Int("add_cpu") + addMemSize, _ := body.Int("add_mem") + hostutils.DelayTaskWithoutReqctx(ctx, guestman.GetGuestManager().HotplugCpuMem, + &guestman.SGuestHotplugCpuMem{sid, addCpuCount, addMemSize}) + return nil, nil +} + func guestReloadDiskSnapshot(ctx context.Context, sid string, body jsonutils.JSONObject) (interface{}, error) { if !guestman.GetGuestManager().IsGuestExist(sid) { return nil, httperrors.NewNotFoundError("Guest %s not found", sid) diff --git a/pkg/hostman/guestman/guesthelper.go b/pkg/hostman/guestman/guesthelper.go index 6323f424ba..5e3cbf8fa9 100644 --- a/pkg/hostman/guestman/guesthelper.go +++ b/pkg/hostman/guestman/guesthelper.go @@ -48,6 +48,12 @@ type SDriverMirror struct { Desc jsonutils.JSONObject } +type SGuestHotplugCpuMem struct { + Sid string + AddCpuCount int64 + AddMemSize int64 +} + type SReloadDisk struct { Sid string Disk storageman.IDisk diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index cd15102474..7a04bebba6 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -190,6 +190,11 @@ func (m *SGuestManager) IsGuestExist(sid string) bool { } } +func (m *SGuestManager) GetGuestById(sid string) *SKVMGuestInstance { + guest, _ := guestManger.Servers[sid] + return guest +} + func (m *SGuestManager) LoadExistingGuests() { files, err := ioutil.ReadDir(m.ServersPath) if err != nil { @@ -650,6 +655,16 @@ func (m *SGuestManager) StartDriveMirror(ctx context.Context, params interface{} return nil, nil } +func (m *SGuestManager) HotplugCpuMem(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { + hotplugParams, ok := params.(*SGuestHotplugCpuMem) + if !ok { + return nil, hostutils.ParamsError + } + guest := guestManger.Servers[hotplugParams.Sid] + NewGuestHotplugCpuMemTask(ctx, guest, int(hotplugParams.AddCpuCount), int(hotplugParams.AddMemSize)).Start() + return nil, nil +} + func (m *SGuestManager) ExitGuestCleanup() { for _, guest := range m.Servers { guest.ExitCleanup(false) diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index 9530929967..a3f253c53a 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -988,3 +988,135 @@ func (task *SGuestOnlineResizeDiskTask) OnResizeSucc(result string) { params.Add(jsonutils.NewInt(task.sizeMB), "disk_size") hostutils.TaskComplete(task.ctx, params) } + +/** + * GuestHotplugCpuMem +**/ + +type SGuestHotplugCpuMemTask struct { + *SKVMGuestInstance + + ctx context.Context + addCpuCount int + addMemSize int + + originalCpuCount int + addedCpuCount int + + memSlotNewIndex *int +} + +func NewGuestHotplugCpuMemTask( + ctx context.Context, s *SKVMGuestInstance, addCpuCount, addMemSize int, +) *SGuestHotplugCpuMemTask { + return &SGuestHotplugCpuMemTask{ + SKVMGuestInstance: s, + ctx: ctx, + addCpuCount: addCpuCount, + addMemSize: addMemSize, + } +} + +// First at all add cpu count, second add mem size +func (task *SGuestHotplugCpuMemTask) Start() { + if task.addCpuCount > 0 { + task.startAddCpu() + } else if task.addMemSize > 0 { + task.startAddMem() + } else { + task.onSucc() + } +} + +func (task *SGuestHotplugCpuMemTask) startAddCpu() { + if task.addCpuCount > 0 { + task.Monitor.GetCpuCount(task.onGetCpuCount) + } + +} + +func (task *SGuestHotplugCpuMemTask) onGetCpuCount(count int) { + task.originalCpuCount = count + task.doAddCpu() +} + +func (task *SGuestHotplugCpuMemTask) doAddCpu() { + if task.addedCpuCount < task.addCpuCount { + task.Monitor.AddCpu(task.originalCpuCount+task.addedCpuCount, task.onAddCpu) + } else { + task.startAddMem() + } +} + +func (task *SGuestHotplugCpuMemTask) onAddCpu(reason string) { + if len(reason) > 0 { + log.Errorln(reason) + task.onFail(reason) + return + } + task.addedCpuCount += 1 + task.doAddCpu() +} + +func (task *SGuestHotplugCpuMemTask) startAddMem() { + if task.addMemSize > 0 { + task.Monitor.GeMemtSlotIndex(task.onGetSlotIndex) + } else { + task.onSucc() + } +} + +// index little then zero indicate guest not take up slot +func (task *SGuestHotplugCpuMemTask) onGetSlotIndex(index int) { + var newIndex int + if index < 0 { + newIndex = 0 + } else { + newIndex = index + 1 + } + task.memSlotNewIndex = &newIndex + params := map[string]string{ + "id": fmt.Sprintf("mem%d", *task.memSlotNewIndex), + "size": fmt.Sprintf("%dM", task.addMemSize), + } + task.Monitor.ObjectAdd("memory-backend-ram", params, task.onAddMemObject) +} + +func (task *SGuestHotplugCpuMemTask) onAddMemObject(reason string) { + if len(reason) > 0 { + log.Errorln(reason) + cb := func(res string) { log.Infof("%s", res) } + task.Monitor.ObjectDel(fmt.Sprintf("mem%d", *task.memSlotNewIndex), cb) + task.onFail(reason) + return + } + params := map[string]interface{}{ + "id": fmt.Sprintf("dimm%d", *task.memSlotNewIndex), + "memdev": fmt.Sprintf("mem%d", *task.memSlotNewIndex), + } + task.Monitor.DeviceAdd("pc-dimm", params, task.onAddMemDevice) +} + +func (task *SGuestHotplugCpuMemTask) onAddMemDevice(reason string) { + if len(reason) > 0 { + log.Errorln(reason) + task.onFail(reason) + return + } + task.onSucc() +} + +func (task *SGuestHotplugCpuMemTask) onFail(reason string) { + body := jsonutils.NewDict() + if task.addedCpuCount < task.addCpuCount { + body.Set("add_cpu_failed", jsonutils.JSONTrue) + body.Set("added_cpu", jsonutils.NewInt(int64(task.addedCpuCount))) + } else if task.memSlotNewIndex != nil { + body.Set("add_mem_failed", jsonutils.JSONTrue) + } + hostutils.TaskFailed2(task.ctx, reason, body) +} + +func (task *SGuestHotplugCpuMemTask) onSucc() { + hostutils.TaskComplete(task.ctx, nil) +} diff --git a/pkg/hostman/guestman/qemu-kvmhelper.go b/pkg/hostman/guestman/qemu-kvmhelper.go index 530c4f364c..4f187c06b5 100644 --- a/pkg/hostman/guestman/qemu-kvmhelper.go +++ b/pkg/hostman/guestman/qemu-kvmhelper.go @@ -404,10 +404,10 @@ func (s *SKVMGuestInstance) generateStartScript(data *jsonutils.JSONDict) (strin cmd += fmt.Sprintf(" -machine %s,accel=%s", s.getMachine(), accel) cmd += " -k en-us" // #cmd += " -g 800x600" - cmd += fmt.Sprintf(" -smp %d", cpu) + cmd += fmt.Sprintf(" -smp %d,maxcpus=128", cpu) cmd += fmt.Sprintf(" -name %s", name) // #cmd += fmt.Sprintf(" -uuid %s", self.desc["uuid"]) - cmd += fmt.Sprintf(" -m %d", mem) + cmd += fmt.Sprintf(" -m %dM,slots=4,maxmem=262144M", mem) if options.HostOptions.HugepagesOption == "native" { cmd += fmt.Sprintf(" -mem-prealloc -mem-path %s", fmt.Sprintf("/dev/hugepages/%s", uuid)) diff --git a/pkg/hostman/hostutils/hostutils.go b/pkg/hostman/hostutils/hostutils.go index 56f5e07157..ccd357b347 100644 --- a/pkg/hostman/hostutils/hostutils.go +++ b/pkg/hostman/hostutils/hostutils.go @@ -51,6 +51,14 @@ func TaskFailed(ctx context.Context, reason string) { } } +func TaskFailed2(ctx context.Context, reason string, params *jsonutils.JSONDict) { + if taskId := ctx.Value(appctx.APP_CONTEXT_KEY_TASK_ID); taskId != nil { + modules.ComputeTasks.TaskFailed3(GetComputeSession(ctx), taskId.(string), reason, params) + } else { + log.Errorf("Reqeuest task failed missing task id, with reason(%s)", reason) + } +} + func TaskComplete(ctx context.Context, params jsonutils.JSONObject) { if taskId := ctx.Value(appctx.APP_CONTEXT_KEY_TASK_ID); taskId != nil { modules.ComputeTasks.TaskComplete(GetComputeSession(ctx), taskId.(string), params) diff --git a/pkg/hostman/monitor/hmp.go b/pkg/hostman/monitor/hmp.go index a48269ac46..c0d6d52878 100644 --- a/pkg/hostman/monitor/hmp.go +++ b/pkg/hostman/monitor/hmp.go @@ -66,18 +66,22 @@ func (m *HmpMonitor) read(r io.Reader) { } else { // remove reader timeout m.connected = true + m.timeout = false m.rwc.SetReadDeadline(time.Time{}) + go m.query() + go m.OnMonitorConnected() } } log.Errorln("Scan over ...") - if err := scanner.Err(); err != nil { - log.Errorln(err) - if m.connected { - m.connected = false - m.OnMonitorDisConnect(err) - } else { - m.OnMonitorTimeout(err) - } + err := scanner.Err() + if err != nil { + log.Infof("HMP Disconnected: %s", err) + } + if m.timeout { + m.OnMonitorTimeout(err) + } else if m.connected { + m.connected = false + m.OnMonitorDisConnect(err) } m.reading = false } @@ -91,6 +95,10 @@ func (m *HmpMonitor) callBack(res string) { m.callbackQueue = m.callbackQueue[1:] m.mutex.Unlock() if cb != nil { + pos := strings.Index(res, "\r\n") + if pos > 0 { + res = res[pos+2:] + } go cb(res) } } @@ -161,8 +169,8 @@ func (m *HmpMonitor) QueryStatus(callback StringCallback) { } func (m *HmpMonitor) parseStatus(callback StringCallback) StringCallback { - return func(res string) { - strs := strings.Split(res, "\r\n") + return func(output string) { + strs := strings.Split(strings.TrimSuffix(output, "\r\n"), "\r\n") for _, str := range strs { if strings.HasPrefix(str, "VM status:") { callback(strings.TrimSpace(str[len("VM status:"):])) @@ -185,8 +193,8 @@ func (m *HmpMonitor) GetVersion(callback StringCallback) { } func (m *HmpMonitor) GetBlocks(callback func(*jsonutils.JSONArray)) { - var cb = func(res string) { - var lines = strings.Split(res, "\r\n") + var cb = func(output string) { + var lines = strings.Split(strings.TrimSuffix(output, "\r\n"), "\r\n") var mergedOutput = []string{} // merge output @@ -250,6 +258,10 @@ func (m *HmpMonitor) DeviceDel(idstr string, callback StringCallback) { m.Query(fmt.Sprintf("device_del %s", idstr), callback) } +func (m *HmpMonitor) ObjectDel(idstr string, callback StringCallback) { + m.Query(fmt.Sprintf("object_del %s", idstr), callback) +} + func (m *HmpMonitor) DriveAdd(bus string, params map[string]string, callback StringCallback) { var paramsKvs = []string{} for k, v := range params { @@ -288,7 +300,7 @@ func (m *HmpMonitor) GetMigrateStatus(callback StringCallback) { log.Infof("Query migrate status: %s", output) var status string - for _, line := range strings.Split(output, "\n") { + for _, line := range strings.Split(strings.TrimSuffix(output, "\r\n"), "\r\n") { if strings.HasPrefix(line, "Migration status") { status = line[strings.LastIndex(line, " ")+1:] break @@ -302,7 +314,7 @@ func (m *HmpMonitor) GetMigrateStatus(callback StringCallback) { func (m *HmpMonitor) GetBlockJobCounts(callback func(jobs int)) { cb := func(output string) { - lines := strings.Split(output, "\n") + lines := strings.Split(strings.TrimSuffix(output, "\r\n"), "\r\n") if lines[0] == "No active jobs" { callback(0) } else { @@ -315,7 +327,7 @@ func (m *HmpMonitor) GetBlockJobCounts(callback func(jobs int)) { func (m *HmpMonitor) GetBlockJobs(callback func(*jsonutils.JSONArray)) { cb := func(output string) { - lines := strings.Split(output, "\n") + lines := strings.Split(strings.TrimSuffix(output, "\r\n"), "\r\n") if lines[0] == "No active jobs" { callback(nil) } else { @@ -383,3 +395,32 @@ func (m *HmpMonitor) ResizeDisk(driveName string, sizeMB int64, callback StringC cmd := fmt.Sprintf("block_resize %s %d", driveName, sizeMB) m.Query(cmd, callback) } + +func (m *HmpMonitor) GetCpuCount(callback func(count int)) { + var cb = func(output string) { + cpus := strings.Split(strings.TrimSuffix(output, "\r\n"), "\r\n") + callback(len(cpus)) + } + m.Query("info cpus", cb) +} + +func (m *HmpMonitor) AddCpu(cpuIndex int, callback StringCallback) { + m.Query(fmt.Sprintf("cpu-add %d", cpuIndex), callback) +} + +func (m *HmpMonitor) GeMemtSlotIndex(callback func(index int)) { + var cb = func(output string) { + memInfos := strings.Split(strings.TrimSuffix(output, "\r\n"), "\r\n") + var count int + for _, line := range memInfos { + if strings.HasPrefix(line, "slot:") { + count += 1 + } + } + if count == 0 { + count = -1 + } + callback(count) + } + m.Query("info memory-devices", cb) +} diff --git a/pkg/hostman/monitor/hmp_test.go b/pkg/hostman/monitor/hmp_test.go index b00ad76a9e..1ab795a358 100644 --- a/pkg/hostman/monitor/hmp_test.go +++ b/pkg/hostman/monitor/hmp_test.go @@ -3,24 +3,25 @@ package monitor import ( "testing" "time" - - "yunion.io/x/log" ) func TestHmpMonitor_Connect(t *testing.T) { - onConnected := func() { log.Infof("Monitor Connected") } - onDisConnect := func(error) { log.Infof("Monitor DisConnect") } - onTimeout := func(error) { log.Infof("Monitor Timeout") } + onConnected := func() { t.Logf("Monitor Connected") } + onDisConnect := func(error) { t.Logf("Monitor DisConnect") } + onTimeout := func(error) { t.Logf("Monitor Timeout") } m := NewHmpMonitor(onDisConnect, onTimeout, onConnected) var host = "127.0.0.1" var port = 55901 m.Connect(host, port) - rawCallBack := func(res string) { log.Infof("OnCallback: %s", res) } + rawCallBack := func(res string) { t.Logf("OnCallback: %s", res) } m.Query("info block", rawCallBack) m.Query("unknown cmd", rawCallBack) - statusCallBack := func(res string) { log.Infof("OnStatusCallback %s", res) } + statusCallBack := func(res string) { t.Logf("OnStatusCallback %s", res) } m.QueryStatus(statusCallBack) - m.Disconnect() + + m.GetBlockJobCounts(func(res int) { t.Logf("GetBlockJobCounts %d", res) }) + m.GetCpuCount(func(res int) { t.Logf("GetCpuCount %d", res) }) time.Sleep(3 * time.Second) + m.Disconnect() } diff --git a/pkg/hostman/monitor/monitor.go b/pkg/hostman/monitor/monitor.go index 2e0151c150..74a5a55efe 100644 --- a/pkg/hostman/monitor/monitor.go +++ b/pkg/hostman/monitor/monitor.go @@ -26,13 +26,19 @@ type Monitor interface { GetBlockJobCounts(func(jobs int)) GetBlockJobs(func(*jsonutils.JSONArray)) + GetCpuCount(func(count int)) + AddCpu(cpuIndex int, callback StringCallback) + GeMemtSlotIndex(func(index int)) + GetBlocks(callback func(*jsonutils.JSONArray)) EjectCdrom(dev string, callback StringCallback) ChangeCdrom(dev string, path string, callback StringCallback) DriveDel(idstr string, callback StringCallback) DeviceDel(idstr string, callback StringCallback) + ObjectDel(idstr string, callback StringCallback) + ObjectAdd(objectType string, params map[string]string, callback StringCallback) DriveAdd(bus string, params map[string]string, callback StringCallback) DeviceAdd(dev string, params map[string]interface{}, callback StringCallback) diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index 01cc6c6561..c20fb02652 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -471,6 +471,10 @@ func (m *QmpMonitor) DeviceDel(idstr string, callback StringCallback) { // m.Query(cmd, cb) } +func (m *QmpMonitor) ObjectDel(idstr string, callback StringCallback) { + m.HumanMonitorCommand(fmt.Sprintf("object_del %s", idstr), callback) +} + func (m *QmpMonitor) DriveAdd(bus string, params map[string]string, callback StringCallback) { var paramsKvs = []string{} for k, v := range params { @@ -723,3 +727,56 @@ func (m *QmpMonitor) ResizeDisk(driveName string, sizeMB int64, callback StringC cmd := fmt.Sprintf("block_resize %s %d", driveName, sizeMB) m.HumanMonitorCommand(cmd, callback) } + +func (m *QmpMonitor) GetCpuCount(callback func(count int)) { + var cb = func(res string) { + cpus := strings.Split(res, "\\n") + count := 0 + for _, cpuInfo := range cpus { + if len(strings.TrimSpace(cpuInfo)) > 0 { + count += 1 + } + } + callback(count) + } + m.HumanMonitorCommand("info cpus", cb) +} + +func (m *QmpMonitor) AddCpu(cpuIndex int, callback StringCallback) { + var ( + cb = func(res *Response) { + callback(m.actionResult(res)) + } + cmd = &Command{ + Execute: "cpu-add", + Args: map[string]interface{}{"id": cpuIndex}, + } + ) + m.Query(cmd, cb) +} + +func (m *QmpMonitor) ObjectAdd(objectType string, params map[string]string, callback StringCallback) { + var paramsKvs = []string{} + for k, v := range params { + paramsKvs = append(paramsKvs, fmt.Sprintf("%s=%s", k, v)) + } + cmd := fmt.Sprintf("object_add %s,%s", objectType, strings.Join(paramsKvs, ",")) + m.HumanMonitorCommand(cmd, callback) +} + +func (m *QmpMonitor) GeMemtSlotIndex(callback func(index int)) { + var cb = func(res string) { + memInfos := strings.Split(res, "\\n") + var count int + for _, line := range memInfos { + if strings.HasPrefix(strings.TrimSpace(line), "slot:") { + count += 1 + } + } + if count == 0 { + count = -1 + } + callback(count) + } + m.HumanMonitorCommand("info memory-devices", cb) +} diff --git a/pkg/mcclient/modules/mod_tasks.go b/pkg/mcclient/modules/mod_tasks.go index 1717a96ca8..cc509f8075 100644 --- a/pkg/mcclient/modules/mod_tasks.go +++ b/pkg/mcclient/modules/mod_tasks.go @@ -49,7 +49,13 @@ func (man ComputeTasksManager) TaskFailed(session *mcclient.ClientSession, taskI } func (man ComputeTasksManager) TaskFailed2(session *mcclient.ClientSession, taskId string, reason string) { - params := jsonutils.NewDict() + man.TaskFailed3(session, taskId, reason, nil) +} + +func (man ComputeTasksManager) TaskFailed3(session *mcclient.ClientSession, taskId string, reason string, params *jsonutils.JSONDict) { + if params == nil { + params = jsonutils.NewDict() + } params.Add(jsonutils.NewString("error"), "__status__") params.Add(jsonutils.NewString(reason), "__reason__") man.TaskComplete(session, taskId, params)