diff --git a/Gopkg.lock b/Gopkg.lock index bcd51fb9d7..6a367673e1 100644 --- a/Gopkg.lock +++ b/Gopkg.lock @@ -1223,11 +1223,11 @@ [[projects]] branch = "master" - digest = "1:49ffc35ec8d3f7789393cd132acd359e8ac1f5d38c7a2b91c484b041840c62c0" + digest = "1:54554b3c72f4fcbd3c27c5137f9acd9d153d7e7150d77cd509af50c75fcb41b9" name = "yunion.io/x/jsonutils" packages = ["."] pruneopts = "UT" - revision = "d1290e94d4753c1748fc7c89f472a523cc0a5c08" + revision = "7079aada4c7e9a37e4c3b183bddb0ef684e038f2" [[projects]] branch = "master" @@ -1242,7 +1242,7 @@ [[projects]] branch = "master" - digest = "1:211ecaed7d1d87e5d216b598c3cd5b7c47977b73a06308e1b39c08fd03c9c417" + digest = "1:201bff9d99a538dd71fd8e55bbbfda8a301fd89676d0668a2cedc574e1aeaa03" name = "yunion.io/x/pkg" packages = [ "gotypes", @@ -1276,7 +1276,7 @@ "utils", ] pruneopts = "UT" - revision = "9d246215b0a167b153bcdbc7d75031ced7138cf8" + revision = "122c7d4ce76b4611b63d86af8e9bb2610ab55c20" [[projects]] branch = "master" diff --git a/cmd/climc/shell/elasticips.go b/cmd/climc/shell/elasticips.go index 054283940c..74e6e0f281 100644 --- a/cmd/climc/shell/elasticips.go +++ b/cmd/climc/shell/elasticips.go @@ -187,4 +187,16 @@ func init() { return nil }) + type EipPurgeOptions struct { + ID string `help:"ID or name of EIP"` + } + R(&EipPurgeOptions{}, "eip-purge", "Purge EIP db records", func(s *mcclient.ClientSession, args *EipPurgeOptions) error { + result, err := modules.Elasticips.PerformAction(s, args.ID, "purge", nil) + if err != nil { + return err + } + printObject(result) + return nil + }) + } diff --git a/cmd/climc/shell/snapshots.go b/cmd/climc/shell/snapshots.go index 3be5a0ad2b..967abb7387 100644 --- a/cmd/climc/shell/snapshots.go +++ b/cmd/climc/shell/snapshots.go @@ -16,6 +16,8 @@ func init() { Share bool `help:"Show shared snapshots"` DiskType string `help: "Filter by disk type" choices:"sys|data"` Provider string `help: "Cloud provider" choices:"Aliyun|VMware|Azure"` + + Manager string `help:"Show snapshots belongs to a specific cloud provider"` } R(&SnapshotsListOptions{}, "snapshot-list", "Show snapshots", func(s *mcclient.ClientSession, args *SnapshotsListOptions) error { params, err := args.BaseListOptions.Params() @@ -41,6 +43,9 @@ func init() { if len(args.Provider) > 0 { params.Add(jsonutils.NewString(args.Provider), "provider") } + if len(args.Manager) > 0 { + params.Add(jsonutils.NewString(args.Manager), "manager") + } result, err := modules.Snapshots.List(s, params) if err != nil { return err @@ -84,4 +89,17 @@ func init() { printObject(result) return nil }) + + type SnapshotPurgeOptions struct { + ID string `help:"ID or name of Snapshot"` + } + R(&SnapshotPurgeOptions{}, "snapshot-purge", "Purge Snapshot db records", func(s *mcclient.ClientSession, args *SnapshotPurgeOptions) error { + result, err := modules.Snapshots.PerformAction(s, args.ID, "purge", nil) + if err != nil { + return err + } + printObject(result) + return nil + }) + } diff --git a/cmd/climc/shell/vpcs.go b/cmd/climc/shell/vpcs.go index 547a775dd0..79517f2880 100644 --- a/cmd/climc/shell/vpcs.go +++ b/cmd/climc/shell/vpcs.go @@ -11,7 +11,7 @@ func init() { type VpcListOptions struct { options.BaseListOptions Region string `help:"ID or Name of region"` - Manager string `help:"Show regions belongs to the cloud provider"` + Manager string `help:"Show vpcs belongs to the cloud provider"` } R(&VpcListOptions{}, "vpc-list", "List VPCs", func(s *mcclient.ClientSession, args *VpcListOptions) error { var params *jsonutils.JSONDict diff --git a/cmd/scheduler/app/server.go b/cmd/scheduler/app/server.go index e3a8f93be8..7a3ff45fc2 100644 --- a/cmd/scheduler/app/server.go +++ b/cmd/scheduler/app/server.go @@ -64,7 +64,7 @@ func Run(s *SchedulerServer) error { debug := o.GetOptions().LogLevel == "debug" - auth.AsyncInit(s.AuthInfo, debug, true, startSched) + auth.AsyncInit(s.AuthInfo, debug, true, "", "", startSched) return startHTTP(s) } diff --git a/pkg/appsrv/appsrv.go b/pkg/appsrv/appsrv.go index 314bfc185f..38c17a7604 100644 --- a/pkg/appsrv/appsrv.go +++ b/pkg/appsrv/appsrv.go @@ -10,83 +10,19 @@ import ( "time" "yunion.io/x/log" - "yunion.io/x/onecloud/pkg/appctx" - "yunion.io/x/onecloud/pkg/proxy" "yunion.io/x/pkg/trace" "yunion.io/x/pkg/utils" + + "yunion.io/x/onecloud/pkg/appctx" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/proxy" + "yunion.io/x/onecloud/pkg/util/httputils" ) -type responseWriterResponse struct { - count int - err error -} - -type responseWriterChannel struct { - backend http.ResponseWriter - bodyChan chan []byte - bodyResp chan responseWriterResponse - statusChan chan int - statusResp chan bool -} - -func newResponseWriterChannel(backend http.ResponseWriter) responseWriterChannel { - return responseWriterChannel{backend: backend, - bodyChan: make(chan []byte), - bodyResp: make(chan responseWriterResponse), - statusChan: make(chan int), - statusResp: make(chan bool)} -} - -func (w *responseWriterChannel) Header() http.Header { - return w.backend.Header() -} - -func (w *responseWriterChannel) Write(bytes []byte) (int, error) { - w.bodyChan <- bytes - v := <-w.bodyResp - return v.count, v.err -} - -func (w *responseWriterChannel) WriteHeader(status int) { - w.statusChan <- status - <-w.statusResp -} - -func (w *responseWriterChannel) wait() { - stop := false - for !stop { - select { - case bytes, more := <-w.bodyChan: - // log.Print("Recive body ", len(bytes), " more ", more) - if more { - c, e := w.backend.Write(bytes) - w.bodyResp <- responseWriterResponse{count: c, err: e} - } else { - stop = true - } - case status, more := <-w.statusChan: - // log.Print("Recive status ", status, " more ", more) - if more { - w.backend.WriteHeader(status) - w.statusResp <- true - } else { - stop = true - } - } - } -} - -func (w *responseWriterChannel) closeChannels() { - close(w.bodyChan) - close(w.bodyResp) - close(w.statusChan) - close(w.statusResp) -} - type Application struct { name string context context.Context - session *WorkerManager + session *SWorkerManager roots map[string]*RadixNode rootLock *sync.Mutex connMax int @@ -106,19 +42,22 @@ const ( DEFAULT_READ_TIMEOUT = 0 DEFAULT_READ_HEADER_TIMEOUT = 10 * time.Second DEFAULT_WRITE_TIMEOUT = 0 + DEFAULT_PROCESS_TIMEOUT = 15 * time.Second ) func NewApplication(name string, connMax int) *Application { app := Application{name: name, context: context.Background(), connMax: connMax, - session: NewWorkerManager("sessionMan", connMax, DEFAULT_BACKLOG), + session: NewWorkerManager("HttpRequestWorkerManager", connMax, DEFAULT_BACKLOG), roots: make(map[string]*RadixNode), rootLock: &sync.Mutex{}, idleTimeout: DEFAULT_IDLE_TIMEOUT, readTimeout: DEFAULT_READ_TIMEOUT, readHeaderTimeout: DEFAULT_READ_HEADER_TIMEOUT, - writeTimeout: DEFAULT_WRITE_TIMEOUT} + writeTimeout: DEFAULT_WRITE_TIMEOUT, + processTimeout: DEFAULT_PROCESS_TIMEOUT, + } app.SetContext(appctx.APP_CONTEXT_KEY_APP, &app) app.SetContext(appctx.APP_CONTEXT_KEY_APPNAME, app.name) @@ -174,9 +113,6 @@ func (app *Application) AddHandler(method string, prefix string, handler func(co func (app *Application) AddHandler2(method string, prefix string, handler func(context.Context, http.ResponseWriter, *http.Request), metadata map[string]interface{}, name string, tags map[string]string) { log.Debugf("%s - %s", method, prefix) segs := SplitPath(prefix) - // for i := len(this.middlewares) - 1; i >= 0; i -= 1 { - // handler = this.middlewares[i](handler) - // } e := app.getRoot(method).Add(segs, newHandlerInfo(method, segs, handler, metadata, name, tags)) if e != nil { log.Fatalf("Fail to register %s %s: %s", method, prefix, e) @@ -253,10 +189,11 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri hand, ok := handler.(*handlerInfo) if ok { fw := newResponseWriterChannel(w) + worker := make(chan *SWorker) errChan := make(chan interface{}) + ctx, cancel := context.WithTimeout(app.context, app.processTimeout) + defer cancel() app.session.Run(func() { - ctx, cancel := context.WithCancel(app.context) - defer cancel() defer fw.closeChannels() if ctx.Err() == nil { ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_REQUEST_ID, rid) @@ -266,25 +203,32 @@ func (app *Application) defaultHandle(w http.ResponseWriter, r *http.Request, ri if hand.metadata != nil { ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_METADATA, hand.metadata) } - span := trace.StartServerTrace(w, r, hand.GetName(params), app.GetName(), hand.GetTags()) - ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_TRACE, span) - hand.handler(ctx, &fw, r) - span.EndTrace() + func() { + span := trace.StartServerTrace(w, r, hand.GetName(params), app.GetName(), hand.GetTags()) + defer span.EndTrace() + ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_TRACE, span) + hand.handler(ctx, &fw, r) + }() } // otherwise, the task has been timeout - }, errChan) - fw.wait() - runerr := WaitChannel(errChan) - if runerr != nil { - http.Error(w, fmt.Sprintf("Internal error: %s", runerr), http.StatusInternalServerError) + }, worker, errChan) + runErr := fw.wait(ctx, worker, errChan) + if runErr != nil { + switch runErr.(type) { + case *httputils.JSONClientError: + je := runErr.(*httputils.JSONClientError) + httperrors.GeneralServerError(w, je) + default: + httperrors.InternalServerError(w, "Internal server error") + } } return hand } else { - log.Printf("Invalid handler for %s", r.URL) - http.Error(w, "Invalid handler", 500) + log.Errorf("Invalid handler for %s", r.URL) + httperrors.InternalServerError(w, "Invalid handler %s", r.URL) } } else if !isCors { - log.Printf("Handler not found") - http.NotFound(w, r) + log.Errorf("Handler not found") + httperrors.NotFoundError(w, "Handler not found") } return nil } @@ -295,6 +239,7 @@ func (app *Application) addDefaultHandler() { app.AddHandler("POST", "/ping", PingHandler) app.AddHandler("GET", "/ping", PingHandler) // app.AddHandler("OPTIONS", "/", CORSHandler) + app.AddHandler("GET", "/worker_stats", WorkerStatsHandler) } func timeoutHandle(h http.Handler) http.HandlerFunc { diff --git a/pkg/appsrv/response.go b/pkg/appsrv/response.go new file mode 100644 index 0000000000..6bf18eea4c --- /dev/null +++ b/pkg/appsrv/response.go @@ -0,0 +1,95 @@ +package appsrv + +import ( + "context" + "net/http" + + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/httperrors" +) + +type responseWriterResponse struct { + count int + err error +} + +type responseWriterChannel struct { + backend http.ResponseWriter + bodyChan chan []byte + bodyResp chan responseWriterResponse + statusChan chan int + statusResp chan bool +} + +func newResponseWriterChannel(backend http.ResponseWriter) responseWriterChannel { + return responseWriterChannel{backend: backend, + bodyChan: make(chan []byte), + bodyResp: make(chan responseWriterResponse), + statusChan: make(chan int), + statusResp: make(chan bool)} +} + +func (w *responseWriterChannel) Header() http.Header { + return w.backend.Header() +} + +func (w *responseWriterChannel) Write(bytes []byte) (int, error) { + w.bodyChan <- bytes + v := <-w.bodyResp + return v.count, v.err +} + +func (w *responseWriterChannel) WriteHeader(status int) { + w.statusChan <- status + <-w.statusResp +} + +func (w *responseWriterChannel) wait(ctx context.Context, workerChan chan *SWorker, errChan chan interface{}) interface{} { + var err interface{} + var worker *SWorker + stop := false + for !stop { + select { + case worker = <-workerChan: + log.Infof("request is being handled by worker %s", worker) + case <-ctx.Done(): + // ctx deadline reached, timeout + if worker != nil { + worker.Detach("timeout") + } + err = httperrors.NewTimeoutError("request process timeout") + stop = true + case e, more := <-errChan: + if more { + err = e + } else { + stop = true + } + case bytes, more := <-w.bodyChan: + // log.Print("Recive body ", len(bytes), " more ", more) + if more { + c, e := w.backend.Write(bytes) + w.bodyResp <- responseWriterResponse{count: c, err: e} + } else { + stop = true + } + case status, more := <-w.statusChan: + // log.Print("Recive status ", status, " more ", more) + if more { + w.backend.WriteHeader(status) + w.statusResp <- true + } else { + stop = true + } + } + } + return err +} + +func (w *responseWriterChannel) closeChannels() { + close(w.bodyChan) + close(w.bodyResp) + close(w.statusChan) + close(w.statusResp) +} diff --git a/pkg/appsrv/workers.go b/pkg/appsrv/workers.go index 675853be3d..45c48b9a6e 100644 --- a/pkg/appsrv/workers.go +++ b/pkg/appsrv/workers.go @@ -1,64 +1,178 @@ package appsrv import ( + "container/list" + "context" + "fmt" + "net/http" "runtime/debug" "sync" + "yunion.io/x/jsonutils" "yunion.io/x/log" ) -type WorkerManager struct { - name string - queue *Ring - workerCount int - backlog int - activeWorker int - workerLock *sync.Mutex - workerId uint64 +const ( + WORKER_STATE_ACTIVE = 0 + WORKER_STATE_DETACH = 1 +) + +var isDebug = false + +func enableDebug() { + isDebug = true } -func NewWorkerManager(name string, workerCount int, backlog int) *WorkerManager { - manager := WorkerManager{name: name, - queue: NewRing(workerCount * backlog), - workerCount: workerCount, - backlog: backlog, - activeWorker: 0, - workerLock: &sync.Mutex{}, - workerId: 0} +var workerManagers []*SWorkerManager + +func init() { + workerManagers = make([]*SWorkerManager, 0) +} + +type SWorker struct { + id uint64 + state int + container *list.Element + manager *SWorkerManager +} + +func newWorker(id uint64, manager *SWorkerManager) *SWorker { + return &SWorker{ + id: id, + state: WORKER_STATE_ACTIVE, + container: nil, + manager: manager, + } +} + +func (worker *SWorker) isDetached() bool { + worker.manager.workerLock.Lock() + defer worker.manager.workerLock.Unlock() + + return worker.state == WORKER_STATE_DETACH +} + +func (worker *SWorker) run() { + for { + if worker.isDetached() { + if isDebug { + log.Debugf("deteched worker %s, no need to pick up new job", worker) + } + break + } + req := worker.manager.queue.Pop() + if req != nil { + task := req.(*sWorkerTask) + if task.worker != nil { + task.worker <- worker + } + if isDebug { + log.Debugf("start exec task on worker %s", worker) + } + execCallback(task) + if isDebug { + log.Debugf("end exec task on worker %s", worker) + } + } else { + if isDebug { + log.Debugf("no more job, exit worker %s", worker) + } + break + } + } + worker.manager.removeWorker(worker) +} + +func (worker *SWorker) Detach(reason string) { + worker.manager.workerLock.Lock() + defer worker.manager.workerLock.Unlock() + + worker.state = WORKER_STATE_DETACH + worker.manager.activeWorker.removeWithLock(worker) + worker.manager.detachedWorker.addWithLock(worker) + + log.Warningf("detach worker %s due to reason %s", worker, reason) +} + +func (worker *SWorker) String() string { + return fmt.Sprintf("#%d(%d)", worker.id, worker.state) +} + +type SWorkerList struct { + list *list.List +} + +func newWorkerList() SWorkerList { + return SWorkerList{ + list: list.New(), + } +} + +func (wl *SWorkerList) addWithLock(worker *SWorker) { + ele := wl.list.PushBack(worker) + worker.container = ele +} + +func (wl *SWorkerList) removeWithLock(worker *SWorker) { + wl.list.Remove(worker.container) + worker.container = nil +} + +func (wl *SWorkerList) size() int { + return wl.list.Len() +} + +type SWorkerManager struct { + name string + queue *Ring + workerCount int + backlog int + activeWorker SWorkerList + detachedWorker SWorkerList + workerLock *sync.Mutex + workerId uint64 +} + +func NewWorkerManager(name string, workerCount int, backlog int) *SWorkerManager { + manager := SWorkerManager{name: name, + queue: NewRing(workerCount * backlog), + workerCount: workerCount, + backlog: backlog, + activeWorker: newWorkerList(), + detachedWorker: newWorkerList(), + workerLock: &sync.Mutex{}, + workerId: 0} + + workerManagers = append(workerManagers, &manager) return &manager } -type workerTask struct { - task func() - err chan interface{} +type sWorkerTask struct { + task func() + worker chan *SWorker + err chan interface{} } -func (wm *WorkerManager) Run(task func(), err chan interface{}) bool { - ret := wm.queue.Push(&workerTask{task: task, err: err}) +func (wm *SWorkerManager) Run(task func(), worker chan *SWorker, err chan interface{}) bool { + ret := wm.queue.Push(&sWorkerTask{task: task, worker: worker, err: err}) if ret { wm.schedule() } return ret } -func (wm *WorkerManager) workerRun(id uint64) { - //log.Println("Start worker", id) - defer wm.decActiveWorker() - req := wm.queue.Pop() - for req != nil { - wm.execCallback(req.(*workerTask)) - req = wm.queue.Pop() - } - //log.Println("End worker", id) -} - -func (wm *WorkerManager) decActiveWorker() { +func (wm *SWorkerManager) removeWorker(worker *SWorker) { wm.workerLock.Lock() defer wm.workerLock.Unlock() - wm.activeWorker -= 1 + + if worker.state == WORKER_STATE_ACTIVE { + wm.activeWorker.removeWithLock(worker) + } else { + wm.detachedWorker.removeWithLock(worker) + } } -func (wm *WorkerManager) execCallback(task *workerTask) { +func execCallback(task *sWorkerTask) { defer func() { if r := recover(); r != nil { log.Errorf("WorkerManager exec callback error: %s", r) @@ -75,16 +189,59 @@ func (wm *WorkerManager) execCallback(task *workerTask) { } } -func (wm *WorkerManager) schedule() { +func (wm *SWorkerManager) schedule() { wm.workerLock.Lock() defer wm.workerLock.Unlock() - if wm.activeWorker < wm.workerCount && wm.queue.Size() > 0 { - wm.activeWorker += 1 + + if wm.activeWorker.size() < wm.workerCount && wm.queue.Size() > 0 { wm.workerId += 1 - go wm.workerRun(wm.workerId) + worker := newWorker(wm.workerId, wm) + wm.activeWorker.addWithLock(worker) + if isDebug { + log.Debugf("no enough worker, add new worker %s", worker) + } + go worker.run() } } +func (wm *SWorkerManager) ActiveWorkerCount() int { + return wm.activeWorker.size() +} + +func (wm *SWorkerManager) DetachedWorkerCount() int { + return wm.detachedWorker.size() +} + +type SWorkerManagerStates struct { + Name string + Backlog int + MaxWorkerCnt int + ActiveWorkerCnt int + DetachWorkerCnt int +} + +func (wm *SWorkerManager) getState() SWorkerManagerStates { + state := SWorkerManagerStates{} + + state.Name = wm.name + state.Backlog = wm.queue.Size() + state.MaxWorkerCnt = wm.workerCount + state.ActiveWorkerCnt = wm.activeWorker.size() + state.DetachWorkerCnt = wm.detachedWorker.size() + + return state +} + +func WorkerStatsHandler(ctx context.Context, w http.ResponseWriter, r *http.Request) { + stats := make([]SWorkerManagerStates, 0) + for i := 0; i < len(workerManagers); i += 1 { + stats = append(stats, workerManagers[i].getState()) + } + result := jsonutils.NewDict() + result.Add(jsonutils.Marshal(&stats), "workers") + fmt.Fprintf(w, result.String()) +} + func WaitChannel(ch chan interface{}) interface{} { var ret interface{} stop := false diff --git a/pkg/appsrv/workers_test.go b/pkg/appsrv/workers_test.go index 69f4877b9f..bc1527ab37 100644 --- a/pkg/appsrv/workers_test.go +++ b/pkg/appsrv/workers_test.go @@ -6,22 +6,22 @@ import ( ) func TestWorkerManager(t *testing.T) { + enableDebug() startTime := time.Now() - end := make(chan int) + // end := make(chan int) wm := NewWorkerManager("testwm", 2, 10) counter := 0 for i := 0; i < 10; i += 1 { wm.Run(func() { counter += 1 time.Sleep(1 * time.Second) - if counter >= i { - end <- 1 - } - }, nil) + }, nil, nil) + } + for wm.ActiveWorkerCount() != 0 { + time.Sleep(time.Second) } - <-end if time.Since(startTime) < 5*time.Second { - t.Error("Increct timing") + t.Error("Incorrect timing") } } @@ -30,7 +30,7 @@ func TestWorkerManagerError(t *testing.T) { err := make(chan interface{}) wm.Run(func() { panic("Panic inside worker") - }, err) + }, nil, err) e := WaitChannel(err) if e == nil { t.Error("Panic not captured") @@ -38,9 +38,10 @@ func TestWorkerManagerError(t *testing.T) { err = make(chan interface{}) wm.Run(func() { time.Sleep(1 * time.Second) - }, err) + }, nil, err) e = WaitChannel(err) if e != nil { t.Error("Should no error") } + } diff --git a/pkg/cloudcommon/db/taskman/handler.go b/pkg/cloudcommon/db/taskman/handler.go index beb9de938b..23dc6dff15 100644 --- a/pkg/cloudcommon/db/taskman/handler.go +++ b/pkg/cloudcommon/db/taskman/handler.go @@ -8,7 +8,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" ) -var taskWorkMan *appsrv.WorkerManager +var taskWorkMan *appsrv.SWorkerManager func init() { taskWorkMan = appsrv.NewWorkerManager("TaskWorkerManager", 4, 100) @@ -22,5 +22,5 @@ func AddTaskHandler(prefix string, app *appsrv.Application) { func runTask(taskId string, data jsonutils.JSONObject) { taskWorkMan.Run(func() { TaskManager.execTask(taskId, data) - }, nil) + }, nil, nil) } diff --git a/pkg/cloudcommon/db/taskman/localtaskworker.go b/pkg/cloudcommon/db/taskman/localtaskworker.go index f60e8b146f..57eaf00fb5 100644 --- a/pkg/cloudcommon/db/taskman/localtaskworker.go +++ b/pkg/cloudcommon/db/taskman/localtaskworker.go @@ -9,7 +9,7 @@ import ( "yunion.io/x/onecloud/pkg/appsrv" ) -var localTaskWorkerMan *appsrv.WorkerManager +var localTaskWorkerMan *appsrv.SWorkerManager func init() { localTaskWorkerMan = appsrv.NewWorkerManager("LocalTaskWorkerManager", 4, 10) @@ -42,5 +42,5 @@ func LocalTaskRun(task ITask, proc func() (jsonutils.JSONObject, error)) { task.ScheduleRun(data) } - }, nil) + }, nil, nil) } diff --git a/pkg/cloudprovider/resources.go b/pkg/cloudprovider/resources.go index f93abdc326..0935ffb38f 100644 --- a/pkg/cloudprovider/resources.go +++ b/pkg/cloudprovider/resources.go @@ -252,6 +252,7 @@ type ICloudSnapshot interface { GetManagerId() string GetSize() int32 GetDiskId() string + GetDiskType() string Delete() error GetRegionId() string } diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index e18887bd0c..f449d7529f 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -4,8 +4,8 @@ import ( "context" "fmt" "net/http" - "regexp" "strconv" + "strings" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -81,10 +81,9 @@ func (self *SKVMGuestDriver) RequestReloadDiskSnapshot(ctx context.Context, gues } func findVNCPort(results string) int { - reg := regexp.MustCompile(`(\d+\.\d+\.\d+\.\d+):([\d]+)`) - finds := reg.FindStringSubmatch(results) - log.Debugf("finds=%s", finds) - port, _ := strconv.Atoi(finds[2]) + vncInfo := strings.Split(results, "\n") + addrParts := strings.Split(vncInfo[1], ":") + port, _ := strconv.Atoi(addrParts[len(addrParts)-1]) return port } diff --git a/pkg/compute/guestdrivers/virtualization.go b/pkg/compute/guestdrivers/virtualization.go index 38158c87fd..643385639a 100644 --- a/pkg/compute/guestdrivers/virtualization.go +++ b/pkg/compute/guestdrivers/virtualization.go @@ -141,12 +141,12 @@ func (self *SVirtualizedGuestDriver) StartGuestSyncstatusTask(guest *models.SGue } func (self *SVirtualizedGuestDriver) RequestStopGuestForDelete(ctx context.Context, guest *models.SGuest, task taskman.ITask) error { - guestStatus, _ := task.GetParams().GetString("guest_status") - if guestStatus == models.VM_RUNNING { - host := guest.GetHost() - if host != nil && host.Enabled && host.HostStatus == models.HOST_ONLINE && !jsonutils.QueryBoolean(task.GetParams(), "purge", false) { - return guest.StartGuestStopTask(ctx, task.GetUserCred(), true, task.GetTaskId()) - } + host := guest.GetHost() + if host != nil && host.Enabled && host.HostStatus == models.HOST_ONLINE { + return guest.StartGuestStopTask(ctx, task.GetUserCred(), true, task.GetTaskId()) + } + if host != nil && !jsonutils.QueryBoolean(task.GetParams(), "purge", false) { + return fmt.Errorf("fail to contact host") } task.ScheduleRun(nil) return nil diff --git a/pkg/compute/models/cloudproviders.go b/pkg/compute/models/cloudproviders.go index 31afc0984f..d3e9b0e5bf 100644 --- a/pkg/compute/models/cloudproviders.go +++ b/pkg/compute/models/cloudproviders.go @@ -96,6 +96,10 @@ func (self *SCloudprovider) getEipCount() int { return ElasticipManager.Query().Equals("manager_id", self.Id).Count() } +func (self *SCloudprovider) getSnapshotCount() int { + return SnapshotManager.Query().Equals("manager_id", self.Id).Count() +} + func (self *SCloudprovider) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { return self.SEnabledStatusStandaloneResourceBase.ValidateUpdateData(ctx, userCred, query, data) } @@ -425,6 +429,7 @@ type SCloudproviderUsage struct { StorageCount int StorageCacheCount int EipCount int + SnapshotCount int } func (usage *SCloudproviderUsage) isEmpty() bool { @@ -443,6 +448,9 @@ func (usage *SCloudproviderUsage) isEmpty() bool { if usage.EipCount > 0 { return false } + if usage.SnapshotCount > 0 { + return false + } return true } @@ -454,6 +462,7 @@ func (self *SCloudprovider) getUsage() *SCloudproviderUsage { usage.StorageCount = self.getStorageCount() usage.StorageCacheCount = self.getStoragecacheCount() usage.EipCount = self.getEipCount() + usage.SnapshotCount = self.getSnapshotCount() return &usage } diff --git a/pkg/compute/models/cloudregions.go b/pkg/compute/models/cloudregions.go index 3d0e2e9a53..7b84f1ca42 100644 --- a/pkg/compute/models/cloudregions.go +++ b/pkg/compute/models/cloudregions.go @@ -57,6 +57,9 @@ func (self *SCloudregion) ValidateDeleteCondition(ctx context.Context) error { if self.GetZoneCount() > 0 || self.GetVpcCount() > 0 { return httperrors.NewNotEmptyError("not empty cloud region") } + if self.Id == "default" { + return httperrors.NewProtectedResourceError("not allow to delete default cloud region") + } return self.SEnabledStatusStandaloneResourceBase.ValidateDeleteCondition(ctx) } diff --git a/pkg/compute/models/elasticips.go b/pkg/compute/models/elasticips.go index 5c5ad97792..69db4e1e1a 100644 --- a/pkg/compute/models/elasticips.go +++ b/pkg/compute/models/elasticips.go @@ -799,3 +799,30 @@ func (manager *SElasticipManager) TotalCount(projectId string, rangeObj db.IStan usage.EIPUsedCount = q3.Count() return usage } + +func (self *SElasticip) AllowPerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return userCred.IsSystemAdmin() +} + +func (self *SElasticip) PerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + err := self.ValidateDeleteCondition(ctx) + if err != nil { + return nil, err + } + provider := self.GetCloudprovider() + if provider != nil { + if provider.Enabled { + return nil, httperrors.NewInvalidStatusError("Cannot purge elastic_ip on enabled cloud provider") + } + } + err = self.RealDelete(ctx, userCred) + return nil, err +} + +func (self *SElasticip) DoPendingDelete(ctx context.Context, userCred mcclient.TokenCredential) { + if self.Mode == EIP_MODE_INSTANCE_PUBLICIP { + self.SVirtualResourceBase.DoPendingDelete(ctx, userCred) + return + } + self.Dissociate(ctx, userCred) +} \ No newline at end of file diff --git a/pkg/compute/models/guestnetworks.go b/pkg/compute/models/guestnetworks.go index 4fe7a290a0..ae110ec521 100644 --- a/pkg/compute/models/guestnetworks.go +++ b/pkg/compute/models/guestnetworks.go @@ -122,7 +122,7 @@ func (manager *SGuestnetworkManager) newGuestNetwork(ctx context.Context, userCr driver = "virtio" } gn.Driver = driver - if bwLimit > 0 { + if bwLimit >= 0 { gn.BwLimit = bwLimit } diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 3b8603ad0b..2d033e8475 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -1018,6 +1018,12 @@ func (self *SGuest) GetCustomizeColumns(ctx context.Context, userCred mcclient.T extra.Add(jsonutils.NewString(timeutils.FullIsoTime(pendingDeletedAt)), "auto_delete_at") } + isGpu := jsonutils.JSONFalse + if self.isGpu() { + isGpu = jsonutils.JSONTrue + } + extra.Add(isGpu, "is_gpu") + return self.moreExtraInfo(extra) } @@ -1087,6 +1093,13 @@ func (self *SGuest) GetExtraDetails(ctx context.Context, userCred mcclient.Token extra.Add(jsonutils.NewString(eip.IpAddr), "eip") extra.Add(jsonutils.NewString(eip.Mode), "eip_mode") } + + isGpu := jsonutils.JSONFalse + if self.isGpu() { + isGpu = jsonutils.JSONTrue + } + extra.Add(isGpu, "is_gpu") + return self.moreExtraInfo(extra) } @@ -1324,6 +1337,10 @@ func (self *SGuest) getAdminSecurityRules() string { } } +func (self *SGuest) isGpu() bool { + return len(self.GetIsolatedDevices()) != 0 +} + func (self *SGuest) GetIsolatedDevices() []SIsolatedDevice { return IsolatedDeviceManager.findAttachedDevicesOfGuest(self) } @@ -1677,6 +1694,10 @@ func (self *SGuest) getMaxDiskIndex() int8 { return int8(len(guestdisks)) } +func (self *SGuest) AttachDisk(disk *SDisk, userCred mcclient.TokenCredential, driver string, cache string, mountpoint string) error { + return self.attach2Disk(disk, userCred, driver, cache, mountpoint) +} + func (self *SGuest) attach2Disk(disk *SDisk, userCred mcclient.TokenCredential, driver string, cache string, mountpoint string) error { if self.isAttach2Disk(disk) { return fmt.Errorf("Guest has been attached to disk") @@ -2449,24 +2470,31 @@ func (self *SGuest) PerformRevokeSecgroup(ctx context.Context, userCred mcclient func (self *SGuest) PerformAssignSecgroup(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { if !utils.IsInStringArray(self.Status, []string{VM_READY, VM_RUNNING, VM_SUSPEND}) { + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, "Cannot assign security rules in status "+self.Status, userCred, false) return nil, httperrors.NewInputParameterError("Cannot assign security rules in status %s", self.Status) } else { if secgrp, err := data.GetString("secgrp"); err != nil { + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, err, userCred, false) return nil, err } else if sg, err := SecurityGroupManager.FetchByIdOrName(userCred.GetProjectId(), secgrp); err != nil { + msg := fmt.Sprintf("SecurityGroup %s not found", secgrp) + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, msg, userCred, false) return nil, httperrors.NewNotFoundError("SecurityGroup %s not found", secgrp) } else { if _, err := self.GetModelManager().TableSpec().Update(self, func() error { self.SecgrpId = sg.GetId() return nil }); err != nil { + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, err, userCred, false) return nil, err } if err := self.StartSyncTask(ctx, userCred, true, ""); err != nil { + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, err, userCred, false) return nil, err } } } + logclient.AddActionLog(self, logclient.ACT_VM_ASSIGNSECGROUP, nil, userCred, true) return nil, nil } @@ -2892,9 +2920,9 @@ func (self *SGuest) PerformChangeBandwidth(ctx context.Context, userCred mcclien return nil, httperrors.NewBadRequestError("Index Not fount or out of NIC index") } bandwidth, err := data.Int("bandwidth") - if err != nil || bandwidth <= 0 { - logclient.AddActionLog(self, logclient.ACT_VM_CHANGE_BANDWIDTH, "Bandwidth must be larger than 0", userCred, false) - return nil, httperrors.NewBadRequestError("Bandwidth must be larger than 0") + if err != nil || bandwidth < 0 { + logclient.AddActionLog(self, logclient.ACT_VM_CHANGE_BANDWIDTH, "Bandwidth must non-negative", userCred, false) + return nil, httperrors.NewBadRequestError("Bandwidth must be non-negative") } guestnic := &guestnics[index] if guestnic.BwLimit != int(bandwidth) { @@ -3084,6 +3112,10 @@ func (self *SGuest) StartChangeConfigTask(ctx context.Context, userCred mcclient } func (self *SGuest) DoPendingDelete(ctx context.Context, userCred mcclient.TokenCredential) { + eip, _ := self.GetEip() + if eip != nil { + eip.DoPendingDelete(ctx, userCred) + } for _, guestdisk := range self.GetDisks() { disk := guestdisk.GetDisk() storage := disk.GetStorage() @@ -4479,10 +4511,18 @@ func (self *SGuest) DeleteEip(ctx context.Context, userCred mcclient.TokenCreden if eip == nil { return nil } - err = eip.Delete(ctx, userCred) - if err != nil { - log.Errorf("Delete eip fail %s", err) - return err + if eip.Mode == EIP_MODE_INSTANCE_PUBLICIP { + err = eip.RealDelete(ctx, userCred) + if err != nil { + log.Errorf("Delete eip on delete server fail %s", err) + return err + } + } else { + err = eip.Dissociate(ctx, userCred) + if err != nil { + log.Errorf("Dissociate eip on delete server fail %s", err) + return err + } } return nil } diff --git a/pkg/compute/models/secgroups.go b/pkg/compute/models/secgroups.go index b407cab27e..c6ac08cbc3 100644 --- a/pkg/compute/models/secgroups.go +++ b/pkg/compute/models/secgroups.go @@ -367,3 +367,14 @@ func (manager *SSecurityGroupManager) InitializeData() error { } return nil } + +func (self *SSecurityGroup) ValidateDeleteCondition(ctx context.Context) error { + cnt := self.GetGuestsCount() + if cnt > 0 { + return httperrors.NewNotEmptyError("the security group is in use") + } + if self.Id == "default" { + return httperrors.NewProtectedResourceError("not allow to delete default security group") + } + return self.SSharableVirtualResourceBase.ValidateDeleteCondition(ctx) +} diff --git a/pkg/compute/models/snapshots.go b/pkg/compute/models/snapshots.go index 3367203ee4..9d2445cd08 100644 --- a/pkg/compute/models/snapshots.go +++ b/pkg/compute/models/snapshots.go @@ -47,6 +47,7 @@ type SSnapshot struct { Size int `nullable:"false" list:"user"` // MB OutOfChain bool `nullable:"false" default:"false" index:"true" list:"admin"` FakeDeleted bool `nullable:"false" default:"false" index:"true"` + DiskType string `width:"32" charset:"ascii" nullable:"true" list:"user"` CloudregionId string `width:"36" charset:"ascii" nullable:"true" list:"user"` } @@ -112,6 +113,18 @@ func (manager *SSnapshotManager) ListItemFilter(ctx context.Context, q *sqlchemy sq := cloudproviderTbl.Query(cloudproviderTbl.Field("id")).Equals("provider", provider) q = q.In("manager_id", sq) } + + if managerStr := jsonutils.GetAnyString(query, []string{"manager", "manager_id"}); len(managerStr) > 0 { + managerObj, err := CloudproviderManager.FetchByIdOrName("", managerStr) + if err != nil { + if err == sql.ErrNoRows { + return nil, httperrors.NewNotFoundError("manager %s not found", managerStr) + } + return nil, httperrors.NewGeneralError(err) + } + q = q.Equals("manager_id", managerObj.GetId()) + } + return q, nil } @@ -132,8 +145,7 @@ func (self *SSnapshot) getMoreDetails(extra *jsonutils.JSONDict) *jsonutils.JSON } disk, _ := self.GetDisk() if disk != nil { - extra.Add(jsonutils.NewString(disk.DiskType), "disk_type") - + // extra.Add(jsonutils.NewString(disk.DiskType), "disk_type") guests := disk.GetGuests() if len(guests) == 1 { extra.Add(jsonutils.NewString(guests[0].Name), "guest") @@ -259,6 +271,7 @@ func (self *SSnapshotManager) CreateSnapshot(ctx context.Context, userCred mccli snapshot.DiskId = disk.Id snapshot.StorageId = disk.StorageId snapshot.Size = disk.DiskSize + snapshot.DiskType = disk.DiskType snapshot.Location = location snapshot.CreatedBy = createdBy snapshot.Name = name @@ -416,10 +429,10 @@ func totalSnapshotCount(projectId string) int { return count } -// Only sync snapshot status func (self *SSnapshot) SyncWithCloudSnapshot(userCred mcclient.TokenCredential, ext cloudprovider.ICloudSnapshot) error { _, err := self.GetModelManager().TableSpec().Update(self, func() error { self.Status = ext.GetStatus() + self.DiskType = ext.GetDiskType() return nil }) if err != nil { @@ -444,6 +457,7 @@ func (manager *SSnapshotManager) newFromCloudSnapshot(userCred mcclient.TokenCre } } + snapshot.DiskType = extSnapshot.GetDiskType() snapshot.Size = int(extSnapshot.GetSize()) * 1024 snapshot.ManagerId = extSnapshot.GetManagerId() snapshot.CloudregionId = region.Id @@ -530,3 +544,22 @@ func (self *SSnapshot) GetISnapshotRegion() (cloudprovider.ICloudRegion, error) } return provider.GetIRegionById(region.GetExternalId()) } + +func (self *SSnapshot) AllowPerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return userCred.IsSystemAdmin() +} + +func (self *SSnapshot) PerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + err := self.ValidateDeleteCondition(ctx) + if err != nil { + return nil, err + } + provider := self.GetCloudprovider() + if provider != nil { + if provider.Enabled { + return nil, httperrors.NewInvalidStatusError("Cannot purge snapshot on enabled cloud provider") + } + } + err = self.RealDelete(ctx, userCred) + return nil, err +} diff --git a/pkg/compute/models/storagecachedimages.go b/pkg/compute/models/storagecachedimages.go index c0e458185c..e42140c4be 100644 --- a/pkg/compute/models/storagecachedimages.go +++ b/pkg/compute/models/storagecachedimages.go @@ -225,7 +225,7 @@ func (self *SStoragecachedimage) isDownloadSessionExpire() bool { } } -func (self *SStoragecachedimage) markDeleting(ctx context.Context, userCred mcclient.TokenCredential) error { +func (self *SStoragecachedimage) markDeleting(ctx context.Context, userCred mcclient.TokenCredential, isForce bool) error { err := self.ValidateDeleteCondition(ctx) if err != nil { return err @@ -237,7 +237,7 @@ func (self *SStoragecachedimage) markDeleting(ctx context.Context, userCred mccl lockman.LockJointObject(ctx, cache, image) defer lockman.ReleaseJointObject(ctx, cache, image) - if utils.IsInStringArray(self.Status, []string{CACHED_IMAGE_STATUS_READY, CACHED_IMAGE_STATUS_DELETING}) { + if !isForce && ! utils.IsInStringArray(self.Status, []string{CACHED_IMAGE_STATUS_READY, CACHED_IMAGE_STATUS_DELETING}) { return httperrors.NewInvalidStatusError("Cannot uncache in status %s", self.Status) } _, err = self.GetModelManager().TableSpec().Update(self, func() error { diff --git a/pkg/compute/models/storagecaches.go b/pkg/compute/models/storagecaches.go index 04546e8429..5e46d08c31 100644 --- a/pkg/compute/models/storagecaches.go +++ b/pkg/compute/models/storagecaches.go @@ -315,7 +315,7 @@ func (self *SStoragecache) PerformUncacheImage(ctx context.Context, userCred mcc return nil, err } - err = scimg.markDeleting(ctx, userCred) + err = scimg.markDeleting(ctx, userCred, isForce) if err != nil { return nil, httperrors.NewInvalidStatusError("Fail to mark cache status: %s", err) } diff --git a/pkg/compute/models/storages.go b/pkg/compute/models/storages.go index 060c266c62..256321c33c 100644 --- a/pkg/compute/models/storages.go +++ b/pkg/compute/models/storages.go @@ -83,7 +83,7 @@ func (manager *SStorageManager) GetContextManager() []db.IModelManager { } func (self *SStorage) ValidateDeleteCondition(ctx context.Context) error { - if self.GetHostCount() > 0 || self.GetDiskCount() > 0 { + if self.GetHostCount() > 0 || self.GetDiskCount() > 0 || self.GetSnapshotCount() > 0 { return httperrors.NewNotEmptyError("Not an empty storage provider") } return self.SEnabledStatusStandaloneResourceBase.ValidateDeleteCondition(ctx) @@ -97,6 +97,10 @@ func (self *SStorage) GetDiskCount() int { return DiskManager.Query().Equals("storage_id", self.Id).Count() } +func (self *SStorage) GetSnapshotCount() int { + return SnapshotManager.Query().Equals("storage_id", self.Id).Count() +} + func (self *SStorage) IsLocal() bool { return self.StorageType == STORAGE_LOCAL || self.StorageType == STORAGE_BAREMETAL } diff --git a/pkg/compute/models/vpcs.go b/pkg/compute/models/vpcs.go index 67fd0cbc82..2ec8e1bd24 100644 --- a/pkg/compute/models/vpcs.go +++ b/pkg/compute/models/vpcs.go @@ -78,6 +78,9 @@ func (self *SVpc) ValidateDeleteCondition(ctx context.Context) error { if self.GetNetworkCount() > 0 { return httperrors.NewNotEmptyError("VPC not empty") } + if self.Id == "default" { + return httperrors.NewProtectedResourceError("not allow to delete default vpc") + } return self.SEnabledStatusStandaloneResourceBase.ValidateDeleteCondition(ctx) } diff --git a/pkg/compute/tasks/guest_delete_task.go b/pkg/compute/tasks/guest_delete_task.go index b54095f3bc..be53871e02 100644 --- a/pkg/compute/tasks/guest_delete_task.go +++ b/pkg/compute/tasks/guest_delete_task.go @@ -26,7 +26,11 @@ func init() { func (self *GuestDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { guest := obj.(*models.SGuest) self.SetStage("on_guest_stop_complete", nil) - guest.GetDriver().RequestStopGuestForDelete(ctx, guest, self) + err := guest.GetDriver().RequestStopGuestForDelete(ctx, guest, self) + if err != nil { + errMsg := jsonutils.NewString(err.Error()) + self.OnGuestStopCompleteFailed(ctx, obj, errMsg) + } } func (self *GuestDeleteTask) OnGuestStopComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { diff --git a/pkg/compute/tasks/guest_detach_disk_task.go b/pkg/compute/tasks/guest_detach_disk_task.go index 0a6653750f..abbe01ea0e 100644 --- a/pkg/compute/tasks/guest_detach_disk_task.go +++ b/pkg/compute/tasks/guest_detach_disk_task.go @@ -26,20 +26,30 @@ func (self *GuestDetachDiskTask) OnInit(ctx context.Context, obj db.IStandaloneM diskId, _ := self.Params.GetString("disk_id") objDisk, err := models.DiskManager.FetchById(diskId) if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, nil, err) return } disk := objDisk.(*models.SDisk) if disk == nil { - self.OnTaskFail(ctx, guest, fmt.Errorf("Connot find disk %s", diskId)) + self.OnTaskFail(ctx, guest, nil, fmt.Errorf("Connot find disk %s", diskId)) return } + guestdisks := disk.GetGuestdisks() + if len(guestdisks) > 0 { + guestdisk := guestdisks[0] + self.Params.Add(jsonutils.NewString(guestdisk.Driver), "driver") + self.Params.Add(jsonutils.NewString(guestdisk.CacheMode), "cache") + self.Params.Add(jsonutils.NewString(guestdisk.Mountpoint), "mountpoint") + } + guest.DetachDisk(ctx, disk, self.UserCred) if disk.Status == models.DISK_INIT { self.OnSyncConfigComplete(ctx, guest, nil) return } + disk.SetStatus(self.UserCred, models.DISK_DETACHING, "Disk detach") + host := guest.GetHost() purge := false if host != nil && host.Status == models.HOST_DISABLED && jsonutils.QueryBoolean(self.Params, "purge", false) { @@ -47,13 +57,12 @@ func (self *GuestDetachDiskTask) OnInit(ctx context.Context, obj db.IStandaloneM } detachStatus, err := guest.GetDriver().GetDetachDiskStatus() if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, disk, err) return } if utils.IsInStringArray(guest.Status, detachStatus) && !purge { self.SetStage("on_sync_config_complete", nil) guest.GetDriver().RequestDetachDisk(ctx, guest, self) - disk.SetStatus(self.UserCred, models.DISK_READY, "Disk detach") } else { self.OnSyncConfigComplete(ctx, guest, nil) } @@ -63,14 +72,15 @@ func (self *GuestDetachDiskTask) OnSyncConfigComplete(ctx context.Context, guest diskId, _ := self.Params.GetString("disk_id") objDisk, err := models.DiskManager.FetchById(diskId) if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, nil, err) return } disk := objDisk.(*models.SDisk) if disk == nil { - self.OnTaskFail(ctx, guest, fmt.Errorf("Connot find disk %s", diskId)) + self.OnTaskFail(ctx, guest, nil, fmt.Errorf("Connot find disk %s", diskId)) return } + disk.SetDiskReady(ctx, self.UserCred, "") keepDisk := jsonutils.QueryBoolean(self.Params, "keep_disk", true) host := guest.GetHost() purge := false @@ -81,12 +91,14 @@ func (self *GuestDetachDiskTask) OnSyncConfigComplete(ctx context.Context, guest db.OpsLog.LogEvent(disk, db.ACT_DELETE, "", self.UserCred) disk.RealDelete(ctx, self.UserCred) self.SetStageComplete(ctx, nil) - } else if (disk.Status == models.DISK_READY || !keepDisk) && disk.GetGuestDiskCount() == 0 && disk.AutoDelete { + return + } + if !keepDisk && disk.GetGuestDiskCount() == 0 && disk.AutoDelete { self.SetStage("on_disk_delete_complete", nil) db.OpsLog.LogEvent(disk, db.ACT_DELETE, "", self.UserCred) err := guest.GetDriver().RequestDeleteDetachedDisk(ctx, disk, self, purge) if err != nil { - self.OnTaskFail(ctx, guest, err) + self.OnTaskFail(ctx, guest, disk, err) return } } else { @@ -94,7 +106,31 @@ func (self *GuestDetachDiskTask) OnSyncConfigComplete(ctx context.Context, guest } } -func (self *GuestDetachDiskTask) OnTaskFail(ctx context.Context, guest *models.SGuest, err error) { +func (self *GuestDetachDiskTask) OnSyncConfigCompleteFailed(ctx context.Context, obj db.IStandaloneModel, resion jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + driver, _ := self.Params.GetString("driver") + cache, _ := self.Params.GetString("cache") + mountpoint, _ := self.Params.GetString("mountpoint") + diskId, _ := self.Params.GetString("disk_id") + objDisk, err := models.DiskManager.FetchById(diskId) + if err != nil { + self.OnTaskFail(ctx, guest, nil, err) + return + } + disk := objDisk.(*models.SDisk) + db.OpsLog.LogEvent(disk, db.ACT_DETACH, resion.String(), self.UserCred) + disk.SetDiskReady(ctx, self.UserCred, "") + err = guest.AttachDisk(disk, self.UserCred, driver, cache, mountpoint) + if err != nil { + self.OnTaskFail(ctx, guest, disk, err) + return + } +} + +func (self *GuestDetachDiskTask) OnTaskFail(ctx context.Context, guest *models.SGuest, disk *models.SDisk, err error) { + if disk != nil { + disk.SetDiskReady(ctx, self.UserCred, "") + } self.SetStageFailed(ctx, err.Error()) log.Errorf("Guest %s GuestDetachDiskTask failed %s", guest.Id, err.Error()) } diff --git a/pkg/httperrors/errors.go b/pkg/httperrors/errors.go index 4e4b226e03..5129f42071 100644 --- a/pkg/httperrors/errors.go +++ b/pkg/httperrors/errors.go @@ -185,6 +185,11 @@ func NewRequireLicenseError(msg string, params ...interface{}) *httputils.JSONCl return NewJsonClientError(402, "RequireLicenseError", msg, err) } +func NewTimeoutError(msg string, params ...interface{}) *httputils.JSONClientError { + msg, err := errorMessage(msg, params...) + return NewJsonClientError(504, "TimeoutError", msg, err) +} + func NewGeneralError(err error) *httputils.JSONClientError { switch err.(type) { case *httputils.JSONClientError: @@ -193,3 +198,8 @@ func NewGeneralError(err error) *httputils.JSONClientError { return NewInternalServerError(err.Error()) } } + +func NewProtectedResourceError(msg string, params ...interface{}) *httputils.JSONClientError { + msg, err := errorMessage(msg, params...) + return NewJsonClientError(403, "ProtectedResourceError(", msg, err) +} diff --git a/pkg/httperrors/httperrors.go b/pkg/httperrors/httperrors.go index b9be4abbaf..aaa6378318 100644 --- a/pkg/httperrors/httperrors.go +++ b/pkg/httperrors/httperrors.go @@ -94,3 +94,11 @@ func TenantNotFoundError(w http.ResponseWriter, msg string, params ...interface{ func OutOfQuotaError(w http.ResponseWriter, msg string, params ...interface{}) { JsonClientError(w, NewOutOfQuotaError(msg, params...)) } + +func TimeoutError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewTimeoutError(msg, params...)) +} + +func ProtectedResourceError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewProtectedResourceError(msg, params...)) +} diff --git a/pkg/util/aliyun/snapshot.go b/pkg/util/aliyun/snapshot.go index 3e896607a0..fc1ee257bc 100644 --- a/pkg/util/aliyun/snapshot.go +++ b/pkg/util/aliyun/snapshot.go @@ -64,6 +64,10 @@ func (self *SSnapshot) GetDiskId() string { return self.SourceDiskId } +func (self *SSnapshot) GetDiskType() string { + return self.SourceDiskType +} + func (self *SSnapshot) Refresh() error { if snapshots, total, err := self.region.GetSnapshots("", "", "", []string{self.SnapshotId}, 0, 1); err != nil { return err diff --git a/pkg/util/azure/snapshot.go b/pkg/util/azure/snapshot.go index 60d6497b22..4b1408d7c1 100644 --- a/pkg/util/azure/snapshot.go +++ b/pkg/util/azure/snapshot.go @@ -172,3 +172,7 @@ func (self *SSnapshot) GetManagerId() string { func (self *SSnapshot) GetRegionId() string { return self.region.GetId() } + +func (self *SSnapshot) GetDiskType() string { + return "" +} diff --git a/pkg/util/logclient/logclient.go b/pkg/util/logclient/logclient.go index 0d8a541788..76796053e3 100644 --- a/pkg/util/logclient/logclient.go +++ b/pkg/util/logclient/logclient.go @@ -55,12 +55,13 @@ const ( ACT_VM_SYNC_CONF = "同步配置" ACT_VM_SYNC_STATUS = "同步状态" ACT_VM_UNBIND_KEYPAIR = "解绑密钥" + ACT_VM_ASSIGNSECGROUP = "关联安全组" ) // golang 不支持 const 的string array, http://t.cn/EzAvbw8 var BLACK_LIST_OBJ_TYPE = []string{"parameter"} -var logclientWorkerMan *appsrv.WorkerManager +var logclientWorkerMan *appsrv.SWorkerManager func init() { logclientWorkerMan = appsrv.NewWorkerManager("LogClientWorkerManager", 1, 50) @@ -119,5 +120,5 @@ func AddActionLog(model IObject, action string, iNotes interface{}, userCred mcc if err != nil { log.Errorf("create action log failed %s", err) } - }, nil) + }, nil, nil) } diff --git a/vendor/yunion.io/x/jsonutils/compond.go b/vendor/yunion.io/x/jsonutils/compond.go new file mode 100644 index 0000000000..1ffc930c35 --- /dev/null +++ b/vendor/yunion.io/x/jsonutils/compond.go @@ -0,0 +1,13 @@ +package jsonutils + +func (val *JSONValue) isCompond() bool { + return false +} + +func (val *JSONDict) isCompond() bool { + return true +} + +func (val *JSONArray) isCompond() bool { + return true +} diff --git a/vendor/yunion.io/x/jsonutils/currency.go b/vendor/yunion.io/x/jsonutils/currency.go new file mode 100644 index 0000000000..0f2d2b77f8 --- /dev/null +++ b/vendor/yunion.io/x/jsonutils/currency.go @@ -0,0 +1,30 @@ +package jsonutils + +import ( + "fmt" + "strings" + "yunion.io/x/pkg/util/regutils" +) + +func normalizeUSCurrency(currency string) string { + return strings.Replace(currency, ",", "", -1) +} + +func normalizeEUCurrency(currency string) string { + commaPos := strings.IndexByte(currency, ',') + if commaPos >= 0 { + return fmt.Sprintf("%s.%s", strings.Replace(currency[:commaPos], ".", "", -1), currency[commaPos+1:]) + } else { + return strings.Replace(currency, ".", "", -1) + } +} + +func normalizeCurrencyString(currency string) string { + if regutils.MatchUSCurrency(currency) { + return normalizeUSCurrency(currency) + } + if regutils.MatchEUCurrency(currency) { + return normalizeEUCurrency(currency) + } + return currency +} diff --git a/vendor/yunion.io/x/jsonutils/interface.go b/vendor/yunion.io/x/jsonutils/interface.go new file mode 100644 index 0000000000..3ee19a10f6 --- /dev/null +++ b/vendor/yunion.io/x/jsonutils/interface.go @@ -0,0 +1,39 @@ +package jsonutils + +func (self *JSONValue) Interface() interface{} { + return nil +} + +func (self *JSONBool) Interface() interface{} { + return self.data +} + +func (self *JSONInt) Interface() interface{} { + return self.data +} + +func (self *JSONFloat) Interface() interface{} { + return self.data +} + +func (self *JSONString) Interface() interface{} { + return self.data +} + +func (self *JSONArray) Interface() interface{} { + ret := make([]interface{}, len(self.data)) + for i := 0; i < len(self.data); i += 1 { + ret[i] = self.data[i].Interface() + } + return ret +} + +func (self *JSONDict) Interface() interface{} { + mapping := make(map[string]interface{}) + + for k, v := range self.data { + mapping[k] = v.Interface() + } + + return mapping +} diff --git a/vendor/yunion.io/x/jsonutils/jsonutils.go b/vendor/yunion.io/x/jsonutils/jsonutils.go index d3321d1caf..2b9f289cfe 100644 --- a/vendor/yunion.io/x/jsonutils/jsonutils.go +++ b/vendor/yunion.io/x/jsonutils/jsonutils.go @@ -65,6 +65,8 @@ type JSONObject interface { Equals(obj JSONObject) bool unmarshalValue(val reflect.Value) error // IsZero() bool + Interface() interface{} + isCompond() bool } type JSONValue struct { diff --git a/vendor/yunion.io/x/jsonutils/unmarshal.go b/vendor/yunion.io/x/jsonutils/unmarshal.go index 79a7a2ac6a..9a1015f3da 100644 --- a/vendor/yunion.io/x/jsonutils/unmarshal.go +++ b/vendor/yunion.io/x/jsonutils/unmarshal.go @@ -321,7 +321,7 @@ func (this *JSONString) unmarshalValue(val reflect.Value) error { } val.SetInt(intVal) case reflect.Float32, reflect.Float64: - floatVal, err := strconv.ParseFloat(this.data, 64) + floatVal, err := strconv.ParseFloat(normalizeCurrencyString(this.data), 64) if err != nil { return err } @@ -348,6 +348,7 @@ func (this *JSONArray) unmarshalValue(val reflect.Value) error { if this.data != nil { array.Add(this.data...) } + val.Set(reflect.ValueOf(array)) return nil case JSONArrayPtrType, JSONObjectType: val.Set(reflect.ValueOf(this)) diff --git a/vendor/yunion.io/x/jsonutils/yamlutils.go b/vendor/yunion.io/x/jsonutils/yamlutils.go index 1daf0a63a3..fd2b978882 100644 --- a/vendor/yunion.io/x/jsonutils/yamlutils.go +++ b/vendor/yunion.io/x/jsonutils/yamlutils.go @@ -61,43 +61,49 @@ func parseYAMLDict(lines []string) (map[string]JSONObject, error) { } else { key := lines[i][0:keypos] val := strings.Trim(lines[i][keypos+1:], " ") + if len(val) > 0 && val != "|" { - o, e := Parse([]byte(val)) - if e != nil { - return dict, e - } else { - dict[key] = o - } + dict[key] = NewString(val) i++ } else { + sublines := make([]string, 0) j := i + 1 for j < len(lines) && len(strings.Trim(lines[j], " ")) == 0 { + sublines = append(sublines, "") j++ } - if j >= len(lines) || lines[j][0] != ' ' { - return dict, fmt.Errorf("Illformat") - } - indent := 0 - for indent < len(lines[j]) && lines[j][indent] == ' ' { - indent++ - } - sublines := make([]string, 0) - for j < len(lines) { - if indent >= len(lines[j]) && len(strings.Trim(lines[j], " ")) == 0 { - j++ - } else if indent < len(lines[j]) && len(strings.Trim(lines[j][:indent], " ")) == 0 { - sublines = append(sublines, lines[j][indent:]) - j++ - } else { - break + if j < len(lines) { + if lines[j][0] != ' ' { + return dict, fmt.Errorf("Illformat") + } + + indent := 0 + for indent < len(lines[j]) && lines[j][indent] == ' ' { + indent++ + } + + for j < len(lines) { + if indent >= len(lines[j]) && len(strings.Trim(lines[j], " ")) == 0 { + sublines = append(sublines, "") + j++ + } else if indent < len(lines[j]) && len(strings.Trim(lines[j][:indent], " ")) == 0 { + sublines = append(sublines, lines[j][indent:]) + j++ + } else { + break + } } } - o, e := parseYAMLLines(sublines) - if e != nil { - return dict, e + if val == "|" { + dict[key] = NewString(strings.Join(sublines, "\n")) } else { + o, e := parseYAMLLines(sublines) + if e != nil { + return dict, e + } dict[key] = o } + i = j } } @@ -192,8 +198,14 @@ func (this *JSONDict) yamlLines() []string { var ret = make([]string, 0) for _, key := range this.SortedKeys() { val := this.data[key] + if val.IsZero() { + switch val.(type) { + case *JSONString, *JSONDict, *JSONArray, *JSONValue: + continue + } + } lines := val.yamlLines() - if len(lines) == 1 { + if !val.isCompond() && len(lines) == 1 { ret = append(ret, fmt.Sprintf("%s: %s", key, lines[0])) } else { switch val.(type) { diff --git a/vendor/yunion.io/x/pkg/util/regutils/regutils.go b/vendor/yunion.io/x/pkg/util/regutils/regutils.go index 5ae4abb988..19fbb24f0f 100644 --- a/vendor/yunion.io/x/pkg/util/regutils/regutils.go +++ b/vendor/yunion.io/x/pkg/util/regutils/regutils.go @@ -32,6 +32,8 @@ var RFC2882_TIME_REG *regexp.Regexp var EMAIL_REG *regexp.Regexp var CHINA_MOBILE_REG *regexp.Regexp var FS_FORMAT_REG *regexp.Regexp +var US_CURRENCY_REG *regexp.Regexp +var EU_CURRENCY_REG *regexp.Regexp func init() { FUNCTION_REG = regexp.MustCompile(`^\w+\(.*\)$`) @@ -62,6 +64,8 @@ func init() { EMAIL_REG = regexp.MustCompile(`^[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,4}$`) CHINA_MOBILE_REG = regexp.MustCompile(`^1[0-9-]{10}$`) FS_FORMAT_REG = regexp.MustCompile(`^(ext|fat|hfs|xfs|swap|ntfs|reiserfs|ufs|btrfs)`) + US_CURRENCY_REG = regexp.MustCompile(`^(\d{0,3}|((\d{1,3},)+\d{3}))(\.\d*)?$`) + EU_CURRENCY_REG = regexp.MustCompile(`^(\d{0,3}|((\d{1,3}\.)+\d{3}))(,\d*)?$`) } func MatchFunction(str string) bool { @@ -179,3 +183,11 @@ func MatchMobile(str string) bool { func MatchFS(str string) bool { return FS_FORMAT_REG.MatchString(str) } + +func MatchUSCurrency(str string) bool { + return US_CURRENCY_REG.MatchString(str) +} + +func MatchEUCurrency(str string) bool { + return EU_CURRENCY_REG.MatchString(str) +}