This commit is contained in:
wanyaoqi
2019-01-22 16:56:47 +08:00
committed by Zexi Li
parent 9b7f468186
commit ca725b7654
9 changed files with 96 additions and 30 deletions
+28 -1
View File
@@ -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 {
+11 -6
View File
@@ -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["<sid>"]
var action = params["<action>"]
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) {
+14 -5
View File
@@ -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 {
@@ -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)
}
}
+5
View File
@@ -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()
+13
View File
@@ -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) {
+1 -2
View File
@@ -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)
}
+1
View File
@@ -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 {