From ca725b7654ac326950a32493353263e3b099a55d Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Mon, 24 Dec 2018 14:44:15 +0800 Subject: [PATCH] temp --- pkg/cloudcommon/workmanager/manager.go | 29 +++++++++++++- pkg/hostman/guestman/guesthandler.go | 17 +++++--- pkg/hostman/guestman/guestman.go | 19 ++++++--- pkg/hostman/guestman/{kvm.go => qemu-kvm.go} | 39 +++++++++++-------- .../{kvmhelper.go => qemu-kvmhelper.go} | 0 pkg/hostman/monitor/monitor.go | 5 +++ pkg/hostman/monitor/qmp.go | 13 +++++++ pkg/hostman/storageman/disklocal.go | 3 +- pkg/hostman/storageman/storagebase.go | 1 + 9 files changed, 96 insertions(+), 30 deletions(-) rename pkg/hostman/guestman/{kvm.go => qemu-kvm.go} (94%) rename pkg/hostman/guestman/{kvmhelper.go => qemu-kvmhelper.go} (100%) diff --git a/pkg/cloudcommon/workmanager/manager.go b/pkg/cloudcommon/workmanager/manager.go index 9d18686653..4375ba8bdf 100644 --- a/pkg/cloudcommon/workmanager/manager.go +++ b/pkg/cloudcommon/workmanager/manager.go @@ -24,13 +24,21 @@ func (w *SWorkManager) done() { atomic.AddInt32(&w.curCount, -1) } +// If delay task is not panic and task func return err is nil +// task complete will be called, otherwise called task failed +// Params is interface for receive any type, task func should do type assert func (w *SWorkManager) DelayTask(ctx context.Context, task DelayTaskFunc, params interface{}) { + if ctx == nil || ctx.Value(APP_CONTEXT_KEY_TASK_ID) == nil { + w.DelayTaskWithoutTaskid(task, params) + return + } + w.add() go func() { defer w.done() defer func() { if r := recover(); r != nil { - log.Errorln("Delay task recover: ", r) + log.Errorln("DelayTask panic: ", r) switch val := r.(type) { case string: httpclients.TaskFailed(ctx, val) @@ -43,13 +51,32 @@ func (w *SWorkManager) DelayTask(ctx context.Context, task DelayTaskFunc, params }() res, err := task(ctx, params) if err != nil { + log.Debugf("DelayTask failed: %s", err) httpclients.TaskFailed(ctx, err.Error()) } else { + log.Debugf("DelayTask complete: %v", res) httpclients.TaskComplete(ctx, res) } }() } +func StartWorker() + +func (w *SWorkManager) DelayTaskWithoutTaskid(task DelayTaskFunc, params interface{}) { + w.add() + go func() { + defer w.done() + defer func() { + if r := recover(); r != nil { + log.Errorln("DelayTaskWithoutTaskid panic: ", r) + } + }() + if _, err := task(ctx, params); err != nil { + log.Errorln("DelayTaskWithoutTaskid", err) + } + }() +} + func (w *SWorkManager) Stop() { log.Infof("WorkManager To stop, wait for workers ...") for w.curCount > 0 { diff --git a/pkg/hostman/guestman/guesthandler.go b/pkg/hostman/guestman/guesthandler.go index 1f0667b0fe..6ca9450d70 100644 --- a/pkg/hostman/guestman/guesthandler.go +++ b/pkg/hostman/guestman/guesthandler.go @@ -24,6 +24,9 @@ func AddGuestTaskHandler(prefix string, app *appsrv.Application) { func guestActions(ctx context.Context, w http.ResponseWriter, r *http.Request) { params, _, body := appsrv.FetchEnv(ctx, w, r) + if body == nil { + body = jsonutils.NewDict() + } var sid = params[""] var action = params[""] if f, ok := actionFuncs[action]; !ok { @@ -87,7 +90,7 @@ func response(ctx context.Context, w http.ResponseWriter, res interface{}) { func doCreate(ctx context.Context, sid string, body jsonutils.JSONObject) (interface{}, error) { err := guestManger.PrepareCreate(sid) if err != nil { - return nil, httperrors.NewBadRequestError(err.Error()) + return nil, err } wm.DelayTask(ctx, guestManger.DoDeploy, &SGuestDeploy{sid, body, true}) return nil, nil @@ -96,20 +99,22 @@ func doCreate(ctx context.Context, sid string, body jsonutils.JSONObject) (inter func doDeploy(ctx context.Context, sid string, body jsonutils.JSONObject) (interface{}, error) { err := guestManger.PrepareDeploy(sid) if err != nil { - return nil, httperrors.NewBadRequestError(err.Error()) + return nil, err } wm.DelayTask(ctx, guestManger.DoDeploy, &SGuestDeploy{sid, body, false}) return nil, nil } func doStart(ctx context.Context, sid string, body jsonutils.JSONObject) (interface{}, error) { - res, err := guestManger.Start(ctx, sid, body) - return nil, nil + return guestManger.GuestStart(ctx, sid, body) } func doStop(ctx context.Context, sid string, body jsonutils.JSONObject) (interface{}, error) { - // TODO - return nil, nil + timeout, err := body.Int("timeout") + if err != nil { + timeout = 30 + } + return nil, guestManger.GuestStop(ctx, sid, timeout) } func doMonitor(ctx context.Context, sid string, body jsonutils.JSONObject) (interface{}, error) { diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index bfac06509e..8e6fc319b5 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -179,7 +179,7 @@ func (m *SGuestManager) PrepareCreate(sid string) error { m.ServersLock.Lock() defer m.ServersLock.Unlock() if _, ok := m.Servers[sid]; ok { - return fmt.Errorf("Guest %s exists", sid) + return httperrors.NewBadRequestError("Guest %s exists", sid) } guest := NewKVMGuestInstance(sid, m) m.Servers[sid] = guest @@ -190,10 +190,10 @@ func (m *SGuestManager) PrepareDeploy(sid string) error { m.ServersLock.Lock() defer m.ServersLock.Unlock() if guest, ok := m.Servers[sid]; !ok { - return fmt.Errorf("Guest %s not exists", sid) + return httperrors.NewBadRequestError("Guest %s not exists", sid) } else { if guest.IsRunning() || guest.IsSuspend() { - return fmt.Errorf("Cannot deploy on running/suspend guest") + return httperrors.NewBadRequestError("Cannot deploy on running/suspend guest") } } return nil @@ -281,7 +281,7 @@ func (m *SGuestManager) Delete(sid string) (*SKVMGuestInstance, error) { } } -func (m *SGuestManager) Start(ctx context.Context, sid string, body jsonutils.JSONObject) (jsonutils.JSONObject, error) { +func (m *SGuestManager) GuestStart(ctx context.Context, sid string, body jsonutils.JSONObject) (jsonutils.JSONObject, error) { if guest, ok := m.Servers[sid]; ok { if desc, err := body.Get("desc"); err != nil { guest.SaveDesc(desc) @@ -301,7 +301,7 @@ func (m *SGuestManager) Start(ctx context.Context, sid string, body jsonutils.JS res.Set("is_running", jsonutils.JSONTrue) return res, nil } else { - return nil, httperrors.NewBadGatewayError("Seems started, but no VNC info") + return nil, httperrors.NewBadRequestError("Seems started, but no VNC info") } } } else { @@ -309,6 +309,15 @@ func (m *SGuestManager) Start(ctx context.Context, sid string, body jsonutils.JS } } +func (m *SGuestManager) GuestStop(ctx context.Context, sid string, timeout int64) error { + if guest, ok := m.Servers[sid]; !ok { + guest.ExecStopTask(ctx, timeout) + return nil + } else { + return httperrors.NewNotFoundError("Guest %s not found", sid) + } +} + func (m *SGuestManager) GetFreeVncPort() int64 { vncPorts := make(map[int]struct{}, 0) for _, guest := range m.Servers { diff --git a/pkg/hostman/guestman/kvm.go b/pkg/hostman/guestman/qemu-kvm.go similarity index 94% rename from pkg/hostman/guestman/kvm.go rename to pkg/hostman/guestman/qemu-kvm.go index 473eefe0ab..9cb21e4cfc 100644 --- a/pkg/hostman/guestman/kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -171,17 +171,18 @@ func (s *SKVMGuestInstance) DirtyServerRequestStart() { } // Delay Process -func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interface{}) { +func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { data, ok := params.(*jsonutils.JSONDict) if !ok { log.Errorln("asyncScriptStart params error") - return + return nil, fmt.Errorf("Unknown params") } // TODO hostinof.instace().clean_deleted_ports time.Sleep(100 * time.Millisecond) - var isStarted, tried, err = false, 0, nil + var isStarted, tried = false, 0 + var err error for !isStarted && tried < MAX_TRY { tried += 1 @@ -207,19 +208,15 @@ func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interfa } } - s.onAsyncScriptStart(ctx, isStarted, err) -} - -func (s *SKVMGuestInstance) onAsyncScriptStart(ctx context.Context, isStarted bool, err error) { + // is on_async_script_start if isStarted { log.Infof("Async start server %s success!", s.GetName()) s.StartMonitor(ctx) + return nil, nil } else { log.Infof("Async start server %s failed: %s!!!", s.GetName(), err) - if ctx != nil { - httpclients.TaskFailed(ctx, fmt.Sprintf("Async start server failed: %s", err)) - } - s.SyncStatus() + cloudcommon.AddTimeout(100*time.Millisecond, s.SyncStatus()) + return nil, err } } @@ -260,7 +257,7 @@ func (s *SKVMGuestInstance) ImportServer(pendingDelete bool) { if s.IsRunning() { log.Infof("%s is running, pending_delete=%s", s.GetName(), pendingDelete) if !pendingDelete { - go s.StartMonitor(nil) + s.StartMonitor(nil) } } else { var action = "stopped" @@ -287,6 +284,10 @@ func (s *SKVMGuestInstance) IsSuspend() bool { return false } +func (s *SKVMGuestInstance) IsMonitorAlive() bool { + return s.monitor != nil && s.monitor.IsConnected() +} + func (s *SKVMGuestInstance) ListStateFilePaths() []string { files, err := ioutil.ReadDir(s.HomeDir()) if err == nil { @@ -301,11 +302,8 @@ func (s *SKVMGuestInstance) ListStateFilePaths() []string { return nil } -// Must called in new goroutine func (s *SKVMGuestInstance) StartMonitor(ctx context.Context) { - // delay 100ms start monitor // cloudcommon.AddTimeout(100*time.Millisecond, func() { s.delayStartMonitor(ctx) }) - time.Sleep(100 * time.Millisecond) - s.delayStartMonitor(ctx) + cloudcommon.AddTimeout(100*time.Millisecond, func() { s.delayStartMonitor(ctx) }) } func (s *SKVMGuestInstance) delayStartMonitor(ctx context.Context) { @@ -504,3 +502,12 @@ func (s *SKVMGuestInstance) Delete(ctx context.Context, migrated bool) error { s.delTmpDisks(ctx, migrated) return exec.Command("rm", "-rf", s.HomeDir()).Run() } + +// ================ Guest Stop Task ================== + +func (s *SKVMGuestInstance) ExecStopTask(ctx context.Context, timeout int64) { + if s.IsRunning() && s.IsMonitorAlive() { + // Do Powerdown + s.monitor.SimpleCommand("system_powerdown", callback) + } +} diff --git a/pkg/hostman/guestman/kvmhelper.go b/pkg/hostman/guestman/qemu-kvmhelper.go similarity index 100% rename from pkg/hostman/guestman/kvmhelper.go rename to pkg/hostman/guestman/qemu-kvmhelper.go diff --git a/pkg/hostman/monitor/monitor.go b/pkg/hostman/monitor/monitor.go index 4fbf068711..1358ae3a64 100644 --- a/pkg/hostman/monitor/monitor.go +++ b/pkg/hostman/monitor/monitor.go @@ -14,6 +14,7 @@ type StringCallback func(string) type Monitor interface { Connect(host string, port int) error Dicconnect() + IsConnected() bool // The callback function will be called in another goroutine SimpleCommand(cmd string, callback StringCallback) @@ -67,6 +68,10 @@ func (m *SBaseMonitor) Disconnect() { } } +func (m *SBaseMonitor) IsConnected() bool { + return m.connected +} + func (m *SBaseMonitor) checkReading() bool { m.mutex.Lock() defer m.mutex.Unlock() diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index 92ce8fc1a8..f8cd756ad3 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -5,6 +5,8 @@ import ( "encoding/json" "fmt" "io" + "regexp" + "strings" "time" "yunion.io/x/log" @@ -234,7 +236,18 @@ func (m *QmpMonitor) Connect(host string, port int) error { return nil } +func (m *QmpMonitor) parseCmd(cmd string) string { + re := regexp.MustCompile(`\s+`) + parts := re.Split(strings.TrimSpace(cmd), -1) + if parts[0] == "info" && len(parts) > 1 { + return "query-" + parts[1] + } else { + return parts[0] + } +} + func (m *QmpMonitor) SimpleCommand(cmd string, callback StringCallback) { + cmd = m.parseCmd(cmd) var cb func(res *Response) if callback != nil { cb = func(res *Response) { diff --git a/pkg/hostman/storageman/disklocal.go b/pkg/hostman/storageman/disklocal.go index dab3ff7841..9f79a0cefe 100644 --- a/pkg/hostman/storageman/disklocal.go +++ b/pkg/hostman/storageman/disklocal.go @@ -65,6 +65,5 @@ func (d *SLocalDisk) Delete() error { print 'delete backing-file:', path self.storage.delete_diskfile(path) */ - d.Storage.RemoveDisk(d) - return nil + return d.Storage.RemoveDisk(d) } diff --git a/pkg/hostman/storageman/storagebase.go b/pkg/hostman/storageman/storagebase.go index 96cb9c921b..679c53d9ed 100644 --- a/pkg/hostman/storageman/storagebase.go +++ b/pkg/hostman/storageman/storagebase.go @@ -14,6 +14,7 @@ type IStorage interface { // Find owner disks first, if not found, call create disk GetDiskById(diskId string) IDisk CreateDisk(diskId string) IDisk + RemoveDisk(IDisk) error } type SBaseStorage struct {