Merge pull request #23 in YUNIONIO/onecloud from ~WANYAOQI/onecloud:feature/wyq/kvm-server-deploy to release/2.1.0

* commit 'f04463b5269157e03739412952ee0c38c2688bac':
  server-stop, server-detach-disk, server-syncstatus, server-monitor, server-purge, fix code, add server-suspend, server-cancel-delete server-change-config, fix disk-resize server-deploy
This commit is contained in:
李泽玺
2018-08-09 17:33:57 +08:00
25 changed files with 1048 additions and 59 deletions
+1 -1
View File
@@ -311,7 +311,7 @@ func init() {
return nil
})
R(&ServerOpsOptions{}, "server-sync", "Sync servers status", func(s *mcclient.ClientSession, args *ServerOpsOptions) error {
R(&ServerOpsOptions{}, "server-sync", "Sync servers configures", func(s *mcclient.ClientSession, args *ServerOpsOptions) error {
ret := modules.Servers.BatchPerformAction(s, args.ID, "sync", nil)
printBatchResults(ret, modules.Servers.GetColumns(s))
return nil
+9 -1
View File
@@ -7,6 +7,7 @@ import (
"github.com/yunionio/jsonutils"
"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/appctx"
"github.com/yunionio/onecloud/pkg/appsrv"
"github.com/yunionio/onecloud/pkg/httperrors"
@@ -201,7 +202,14 @@ func createInContextHandler(ctx context.Context, w http.ResponseWriter, r *http.
func performClassActionHandler(ctx context.Context, w http.ResponseWriter, r *http.Request) {
manager, params, query, body := fetchEnv(ctx, w, r)
data, _ := body.Get(manager.KeywordPlural())
var data jsonutils.JSONObject
if body != nil {
data, _ = body.Get(manager.KeywordPlural())
// about string ??
if data == nil {
data = body.(*jsonutils.JSONDict)
}
}
if data == nil {
data = jsonutils.NewDict()
}
+2 -2
View File
@@ -1036,13 +1036,13 @@ func deleteItem(manager IModelManager, model IModel, ctx context.Context, userCr
err := model.ValidateDeleteCondition(ctx)
if err != nil {
log.Errorf("validate delete condition error: %s", err)
return nil, httperrors.NewGeneralError(err)
return nil, httperrors.NewNotAcceptableError(err.Error())
}
err = model.CustomizeDelete(ctx, userCred, query, data)
if err != nil {
log.Errorf("customize delete error: %s", err)
return nil, httperrors.NewGeneralError(err)
return nil, httperrors.NewNotAcceptableError(err.Error())
}
details, err := getItemDetails(manager, model, ctx, userCred, query)
+3 -1
View File
@@ -63,7 +63,9 @@ func (manager *SQuotaManager) _cancelPendingUsage(ctx context.Context, userCred
log.Errorf("%s", err)
return err
}
localUsage.Sub(cancelUsage)
if localUsage != nil {
localUsage.Sub(cancelUsage)
}
quota.Sub(cancelUsage)
err = manager.pendingStore.SetQuota(ctx, userCred, projectId, quota)
if err != nil {
+4 -4
View File
@@ -10,17 +10,17 @@ import (
"github.com/yunionio/jsonutils"
"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/appctx"
"github.com/yunionio/onecloud/pkg/mcclient"
"github.com/yunionio/onecloud/pkg/util/httputils"
"github.com/yunionio/pkg/util/reflectutils"
"github.com/yunionio/pkg/util/stringutils"
"github.com/yunionio/pkg/utils"
"github.com/yunionio/sqlchemy"
"github.com/yunionio/onecloud/pkg/appctx"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/lockman"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/quotas"
"github.com/yunionio/onecloud/pkg/mcclient"
"github.com/yunionio/onecloud/pkg/util/httputils"
)
const (
@@ -387,7 +387,7 @@ func execITask(taskValue reflect.Value, task *STask, data jsonutils.JSONObject,
params[2] = reflect.ValueOf(data)
log.Debugf("Call %s with %s", funcValue, params)
log.Debugf("Call %s: %s with %s", stageName, funcValue, params)
funcValue.Call(params)
+2 -2
View File
@@ -8,13 +8,13 @@ import (
"github.com/yunionio/jsonutils"
"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/httperrors"
"github.com/yunionio/onecloud/pkg/mcclient"
"github.com/yunionio/pkg/util/timeutils"
"github.com/yunionio/pkg/utils"
"github.com/yunionio/sqlchemy"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/lockman"
"github.com/yunionio/onecloud/pkg/httperrors"
"github.com/yunionio/onecloud/pkg/mcclient"
)
type SVirtualResourceBaseManager struct {
+14 -1
View File
@@ -2,13 +2,14 @@ package guestdrivers
import (
"context"
"fmt"
"github.com/yunionio/jsonutils"
"github.com/yunionio/onecloud/pkg/mcclient"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/quotas"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman"
"github.com/yunionio/onecloud/pkg/compute/models"
"github.com/yunionio/onecloud/pkg/mcclient"
)
type SBaremetalGuestDriver struct {
@@ -155,7 +156,19 @@ func (self *SBaremetalGuestDriver) RequestDeployGuestOnHost(ctx context.Context,
return nil
}
func (self *SBaremetalGuestDriver) CanKeepDetachDisk() bool {
return false
}
func (self *SBaremetalGuestDriver) RequestSyncConfigOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
task.ScheduleRun(nil)
return nil
}
func (self *SBaremetalGuestDriver) StartGuestDetachdiskTask(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, params *jsonutils.JSONDict, parentTaskId string) error {
return fmt.Errorf("Cannot detach disk from a baremetal serer")
}
func (self *SBaremetalGuestDriver) StartSuspendTask(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, params *jsonutils.JSONDict, parentTaskId string) error {
return fmt.Errorf("Cannot suspend a baremetal serer")
}
+34 -1
View File
@@ -63,7 +63,7 @@ func (self *SBaseGuestDriver) OnGuestCreateTaskComplete(ctx context.Context, gue
//}
}
func (self *SBaseGuestDriver) StartDeleteGuestTask(guest *models.SGuest, ctx context.Context, userCred mcclient.TokenCredential, params *jsonutils.JSONDict, parentTaskId string) error {
func (self *SBaseGuestDriver) StartDeleteGuestTask(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, params *jsonutils.JSONDict, parentTaskId string) error {
task, err := taskman.TaskManager.NewTask(ctx, "GuestDeleteTask", guest, userCred, params, parentTaskId, "", nil)
if err != nil {
return err
@@ -80,3 +80,36 @@ func (self *SBaseGuestDriver) RequestDetachDisksFromGuestForDelete(ctx context.C
func (self *SBaseGuestDriver) OnDeleteGuestFinalCleanup(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential) error {
return guest.DeleteAllDisksInDB(ctx, userCred)
}
func (self *SBaseGuestDriver) RequestDetachDisk(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
task.ScheduleRun(nil)
return nil
}
func (self *SBaseGuestDriver) RequestGuestCreateAllDisks(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
return fmt.Errorf("Not Implement")
}
func (self *SBaseGuestDriver) GetDetachDiskStatus() ([]string, error) {
return []string{}, fmt.Errorf("This Guest driver dose not implement GetDetachDiskStatus")
}
func (self *SBaseGuestDriver) RequestDeleteDetachedDisk(ctx context.Context, disk *models.SDisk, task taskman.ITask, isPurge bool) error {
return fmt.Errorf("Not Implement")
}
func (self *SBaseGuestDriver) RqeuestSuspendOnHost(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
return fmt.Errorf("Not Implement")
}
func (self *SBaseGuestDriver) AllowReconfigGuest() bool {
return true
}
func (self *SBaseGuestDriver) DoGuestCreateDisksTask(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
return fmt.Errorf("Not Implement")
}
func (self *SBaseGuestDriver) RequestChangeVmConfig(ctx context.Context, guest *models.SGuest, task taskman.ITask, vcpuCount, vmemSize int64) error {
return fmt.Errorf("Not Implement")
}
+17
View File
@@ -27,3 +27,20 @@ func (self *SESXiGuestDriver) RequestSyncConfigOnHost(ctx context.Context, guest
task.ScheduleRun(nil)
return nil
}
func (self *SESXiGuestDriver) GetDetachDiskStatus() ([]string, error) {
return []string{models.VM_READY}, nil
}
func (self *SESXiGuestDriver) CanKeepDetachDisk() bool {
return false
}
func (self *SESXiGuestDriver) RequestDeleteDetachedDisk(ctx context.Context, disk *models.SDisk, task taskman.ITask, isPurge bool) error {
err := disk.RealDelete(ctx, task.GetUserCred())
if err != nil {
return err
}
task.ScheduleRun(nil)
return nil
}
+71 -19
View File
@@ -38,9 +38,12 @@ func (self *SKVMGuestDriver) RequestDetachDisksFromGuestForDelete(ctx context.Co
return nil
}
func (self *SKVMGuestDriver) OnDeleteGuestFinalCleanup(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential) error {
// guest.DeleteAllDisksInDB(ctx, userCred)
// do nothing
func (self *SKVMGuestDriver) DoGuestCreateDisksTask(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
subtask, err := taskman.TaskManager.NewTask(ctx, "KVMGuestCreateDiskTask", guest, task.GetUserCred(), task.GetParams(), task.GetTaskId(), "", nil)
if err != nil {
return err
}
subtask.ScheduleRun(nil)
return nil
}
@@ -112,6 +115,43 @@ func (self *SKVMGuestDriver) RequestStopOnHost(ctx context.Context, guest *model
url := fmt.Sprintf("%s/servers/%s/stop", host.ManagerUri, guest.Id)
_, _, err = httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, body, false)
return err
}
func (self *SKVMGuestDriver) RequestUndeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
url := fmt.Sprintf("%s/servers/%s", host.ManagerUri, guest.Id)
header := http.Header{}
header.Set("X-Auth-Token", task.GetUserCred().GetTokenString())
header.Set("X-Task-Id", task.GetTaskId())
header.Set("X-Region-Version", "v2")
_, res, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "DELETE", url, header, nil, false)
if err != nil {
return err
}
delayClean := jsonutils.QueryBoolean(res, "delay_clean", false)
if res != nil && delayClean {
return nil
}
task.ScheduleRun(nil)
return nil
}
func (self *SKVMGuestDriver) RequestDeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
config := guest.GetDeployConfigOnHost(ctx, host, task.GetParams())
log.Debugf("RequestDeployGuestOnHost: %s", config)
if config.Contains("container") {
// ...
}
action, err := config.GetString("action")
if err != nil {
return err
}
url := fmt.Sprintf("%s/servers/%s/%s", host.ManagerUri, guest.Id, action)
header := http.Header{}
header.Set("X-Auth-Token", task.GetUserCred().GetTokenString())
header.Set("X-Task-Id", task.GetTaskId())
header.Set("X-Region-Version", "v2")
_, _, err = httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, config, false)
if err != nil {
return err
}
@@ -119,13 +159,8 @@ func (self *SKVMGuestDriver) RequestStopOnHost(ctx context.Context, guest *model
}
func (self *SKVMGuestDriver) OnGuestDeployTaskDataReceived(ctx context.Context, guest *models.SGuest, task taskman.ITask, data jsonutils.JSONObject) error {
// ToDO
return fmt.Errorf("Not Implement")
}
func (self *SKVMGuestDriver) RequestUndeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
// ToDo
return fmt.Errorf("Not Implement")
guest.SaveDeployInfo(ctx, task.GetUserCred(), data)
return nil
}
func (self *SKVMGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) (jsonutils.JSONObject, error) {
@@ -162,19 +197,25 @@ func (self *SKVMGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest
return res, nil
}
func (self *SKVMGuestDriver) RequestDeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
// ToDo
return fmt.Errorf("Not Implement")
func (self *SKVMGuestDriver) OnDeleteGuestFinalCleanup(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential) error {
return nil
}
func (self *SKVMGuestDriver) RequestGuestCreateInsertIso(ctx context.Context, imageId string, guest *models.SGuest, task taskman.ITask) error {
// ToDo
return fmt.Errorf("Not Implement")
func (self *SKVMGuestDriver) RequestChangeVmConfig(ctx context.Context, guest *models.SGuest, task taskman.ITask, vcpuCount, vmemSize int64) error {
// pass
return nil
}
func (self *SKVMGuestDriver) RequestGuestCreateAllDisks(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
// ToDo
return fmt.Errorf("Not Implement")
func (self *SKVMGuestDriver) RequestDetachDisk(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
return guest.StartSyncTask(ctx, task.GetUserCred(), false, task.GetTaskId())
}
func (self *SKVMGuestDriver) GetDetachDiskStatus() ([]string, error) {
return []string{models.VM_READY, models.VM_RUNNING}, nil
}
func (self *SKVMGuestDriver) RequestDeleteDetachedDisk(ctx context.Context, disk *models.SDisk, task taskman.ITask, isPurge bool) error {
return disk.StartDiskDeleteTask(ctx, task.GetUserCred(), task.GetTaskId(), isPurge)
}
func (self *SKVMGuestDriver) RequestSyncConfigOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
@@ -191,3 +232,14 @@ func (self *SKVMGuestDriver) RequestSyncConfigOnHost(ctx context.Context, guest
_, err := host.Request(task.GetUserCred(), "POST", url, header, body)
return err
}
func (self *SKVMGuestDriver) RqeuestSuspendOnHost(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
host := guest.GetHost()
url := fmt.Sprintf("%s/servers/%s/suspend", host.ManagerUri, guest.Id)
header := http.Header{}
header.Add("X-Auth-Token", task.GetUserCred().GetTokenString())
header.Add("X-Task-Id", task.GetTaskId())
header.Add("X-Region-Version", "v2")
_, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, nil, false)
return err
}
@@ -89,6 +89,7 @@ func (self *SVirtualizedGuestDriver) StartGuestStopTask(guest *models.SGuest, ct
func (self *SVirtualizedGuestDriver) OnGuestDeployTaskComplete(ctx context.Context, guest *models.SGuest, task taskman.ITask) error {
if jsonutils.QueryBoolean(task.GetParams(), "restart", false) {
task.SetStage("OnDeployStartGuestComplete", nil)
return guest.StartGueststartTask(ctx, task.GetUserCred(), nil, task.GetTaskId())
} else {
guest.SetStatus(task.GetUserCred(), models.VM_READY, "ready")
@@ -149,3 +150,25 @@ func (self *SVirtualizedGuestDriver) CheckDiskTemplateOnStorage(ctx context.Cont
}
return cache.StartImageCacheTask(ctx, userCred, imageId, false, task.GetTaskId())
}
func (self *SVirtualizedGuestDriver) CanKeepDetachDisk() bool {
return true
}
func (self *SVirtualizedGuestDriver) StartGuestDetachdiskTask(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, params *jsonutils.JSONDict, parentTaskId string) error {
task, err := taskman.TaskManager.NewTask(ctx, "GuestDetachDiskTask", guest, userCred, params, parentTaskId, "", nil)
if err != nil {
return err
}
task.ScheduleRun(nil)
return nil
}
func (self *SVirtualizedGuestDriver) StartSuspendTask(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, params *jsonutils.JSONDict, parentTaskId string) error {
task, err := taskman.TaskManager.NewTask(ctx, "GuestSuspendTask", guest, userCred, params, parentTaskId, "", nil)
if err != nil {
return err
}
task.ScheduleRun(nil)
return nil
}
+9 -6
View File
@@ -9,8 +9,6 @@ import (
"github.com/yunionio/jsonutils"
"github.com/yunionio/log"
"github.com/yunionio/sqlchemy"
"github.com/yunionio/pkg/tristate"
"github.com/yunionio/pkg/util/compare"
"github.com/yunionio/pkg/util/fileutils"
@@ -18,6 +16,7 @@ import (
"github.com/yunionio/pkg/util/regutils"
"github.com/yunionio/pkg/util/sysutils"
"github.com/yunionio/pkg/utils"
"github.com/yunionio/sqlchemy"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/quotas"
@@ -38,6 +37,7 @@ const (
DISK_DEALLOC = "deallocating"
DISK_DEALLOC_FAILED = "dealloc_failed"
DISK_UNKNOWN = "unknown"
DISK_DETACHING = "detaching"
DISK_START_SAVE = "start_save"
DISK_SAVING = "saving"
@@ -162,11 +162,14 @@ func (self *SDisk) GetGuestdisks() []SGuestdisk {
}
func (self *SDisk) GetGuests() []SGuest {
result := make([]SGuest, 0)
guests := GuestManager.Query().SubQuery()
query := GuestManager.Query()
guestdisks := GuestdiskManager.Query().SubQuery()
if err := guests.Query().Join(guestdisks, sqlchemy.AND(
sqlchemy.Equals(guestdisks.Field("guest_id"), guests.Field("id")))).
Filter(sqlchemy.Equals(guestdisks.Field("disk_id"), self.Id)).All(&result); err != nil {
q := query.Join(guestdisks, sqlchemy.AND(
sqlchemy.Equals(guestdisks.Field("guest_id"), query.Field("id")))).
Filter(sqlchemy.Equals(guestdisks.Field("disk_id"), self.Id))
// q.DebugQuery()
err := db.FetchModelObjects(GuestManager, q, &result)
if err != nil {
log.Errorf(err.Error())
return nil
}
+15 -1
View File
@@ -56,7 +56,7 @@ type IGuestDriver interface {
RequestStopOnHost(ctx context.Context, guest *SGuest, host *SHost, task taskman.ITask) error
StartDeleteGuestTask(guest *SGuest, ctx context.Context, userCred mcclient.TokenCredential, params *jsonutils.JSONDict, parentTaskId string) error
StartDeleteGuestTask(ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, params *jsonutils.JSONDict, parentTaskId string) error
RequestStopGuestForDelete(ctx context.Context, guest *SGuest, task taskman.ITask) error
@@ -71,6 +71,20 @@ type IGuestDriver interface {
CheckDiskTemplateOnStorage(ctx context.Context, userCred mcclient.TokenCredential, imageId string, storageId string, task taskman.ITask) error
GetGuestVncInfo(userCred mcclient.TokenCredential, guest *SGuest, host *SHost) (*jsonutils.JSONDict, error)
RequestDetachDisk(ctx context.Context, guest *SGuest, task taskman.ITask) error
GetDetachDiskStatus() ([]string, error)
CanKeepDetachDisk() bool
RequestDeleteDetachedDisk(ctx context.Context, disk *SDisk, task taskman.ITask, isPurge bool) error
StartGuestDetachdiskTask(ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, params *jsonutils.JSONDict, parentTaskId string) error
StartSuspendTask(ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, params *jsonutils.JSONDict, parentTaskId string) error
RqeuestSuspendOnHost(ctx context.Context, guest *SGuest, task taskman.ITask) error
AllowReconfigGuest() bool
DoGuestCreateDisksTask(ctx context.Context, guest *SGuest, task taskman.ITask) error
RequestChangeVmConfig(ctx context.Context, guest *SGuest, task taskman.ITask, vcpuCount, vmemSize int64) error
}
var guestDrivers map[string]IGuestDriver
+2
View File
@@ -315,9 +315,11 @@ func (manager *SGuestnetworkManager) DeleteGuestNics(ctx context.Context, guest
if regutils.MatchIP4Addr(gn.IpAddr) || regutils.MatchIP6Addr(gn.Ip6Addr) {
net.updateDnsRecord(&gn, false)
if regutils.MatchIP4Addr(gn.IpAddr) {
// ??
// netman.get_manager().netmap_remove_node(gn.ip_addr)
}
}
// ??
// gn.Delete(ctx, userCred)
err = gn.Delete(ctx, userCred)
if err != nil {
+346 -9
View File
@@ -5,6 +5,7 @@ import (
"context"
"database/sql"
"fmt"
"net/http"
"strconv"
"strings"
"time"
@@ -32,6 +33,7 @@ import (
"github.com/yunionio/onecloud/pkg/httperrors"
"github.com/yunionio/onecloud/pkg/mcclient"
"github.com/yunionio/onecloud/pkg/mcclient/auth"
"github.com/yunionio/onecloud/pkg/util/httputils"
)
const (
@@ -1429,7 +1431,51 @@ func (self *SGuest) PerformSync(ctx context.Context, userCred mcclient.TokenCred
return nil, nil
}
func (self *SGuest) AllowPerformAttachDisk(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
func (self *SGuest) AllowPerformDeploy(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return self.IsOwner(userCred)
}
func (self *SGuest) PerformDeploy(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
kwargs, ok := data.(*jsonutils.JSONDict)
if !ok {
return nil, fmt.Errorf("Parse query body error")
}
if kwargs.Contains("__delete_keypair__") || kwargs.Contains("keypair") {
var kpId string
if !jsonutils.QueryBoolean(kwargs, "__delete_keypair__", false) {
keypair, _ := kwargs.GetString("keypair")
iKp, err := KeypairManager.FetchByIdOrName(userCred.GetProjectId(), keypair)
if err != nil {
return nil, err
}
if iKp == nil {
return nil, fmt.Errorf("Fetch keypair error")
}
kp := iKp.(*SKeypair)
kpId = kp.Id
}
if self.KeypairId != kpId {
self.GetModelManager().TableSpec().Update(self, func() error {
self.KeypairId = kpId
return nil
})
kwargs.Set("reset_password", jsonutils.JSONTrue)
}
}
if utils.IsInStringArray(self.Status, []string{VM_RUNNING, VM_READY, VM_ADMIN}) {
if self.Status == VM_RUNNING {
kwargs.Set("restart", jsonutils.JSONTrue)
}
err := self.StartGuestDeployTask(ctx, userCred, kwargs, "deploy", "")
if err != nil {
return nil, err
}
return nil, nil
}
return nil, httperrors.NewServerStatusError("Cannot deploy in status %s", self.Status)
}
func (self *SGuest) AllowPerformAttachdisk(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return self.IsOwner(userCred)
}
@@ -1446,7 +1492,7 @@ func (self *SGuest) ValidateAttachDisk(ctx context.Context, disk *SDisk) error {
return nil
}
func (self *SGuest) PerformAttachDisk(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
func (self *SGuest) PerformAttachdisk(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
if diskId, err := data.GetString("disk_id"); err != nil {
return nil, err
} else {
@@ -1986,6 +2032,15 @@ func (self *SGuest) StartGueststartTask(ctx context.Context, userCred mcclient.T
return nil
}
func (self *SGuest) StartGuestCreateDiskTask(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict, parentTaskId string) error {
task, err := taskman.TaskManager.NewTask(ctx, "GuestCreateDiskTask", self, userCred, data, parentTaskId, "", nil)
if err != nil {
return err
}
task.ScheduleRun(nil)
return nil
}
func (self *SGuest) StartSyncstatus(ctx context.Context, userCred mcclient.TokenCredential, parentTaskId string) error {
return self.GetDriver().StartGuestSyncstatusTask(self, ctx, userCred, parentTaskId)
}
@@ -2005,7 +2060,7 @@ func (self *SGuest) StartDeleteGuestTask(ctx context.Context, userCred mcclient.
params.Add(jsonutils.JSONTrue, "override_pending_delete")
}
self.SetStatus(userCred, VM_START_DELETE, "")
return self.GetDriver().StartDeleteGuestTask(self, ctx, userCred, params, parentTaskId)
return self.GetDriver().StartDeleteGuestTask(ctx, userCred, self, params, parentTaskId)
}
func (self *SGuest) AllowPerformPurge(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
@@ -2025,13 +2080,206 @@ func (self *SGuest) PerformPurge(ctx context.Context, userCred mcclient.TokenCre
return nil, err
}
func (self *SGuest) detachDisk(ctx context.Context, disk *SDisk, userCred mcclient.TokenCredential) {
func (self *SGuest) DetachDisk(ctx context.Context, disk *SDisk, userCred mcclient.TokenCredential) {
guestdisk := self.GetGuestDisk(disk.Id)
if guestdisk != nil {
guestdisk.Detach(ctx, userCred)
}
}
func (self *SGuest) AllowPerformDetachdisk(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return self.IsOwner(userCred)
}
func (self *SGuest) PerformDetachdisk(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
diskId, err := data.GetString("disk_id")
if err != nil {
return nil, err
}
keepDisk := jsonutils.QueryBoolean(data, "keep_disk", false)
iDisk, err := DiskManager.FetchByIdOrName(userCred.GetProjectId(), diskId)
if err != nil {
return nil, err
}
disk := iDisk.(*SDisk)
if disk != nil {
if self.isAttach2Disk(disk) {
detachDiskStatus, err := self.GetDriver().GetDetachDiskStatus()
if err != nil {
return nil, err
}
if keepDisk && !self.GetDriver().CanKeepDetachDisk() {
return nil, httperrors.NewInputParameterError("Cannot keep detached disk")
}
if utils.IsInStringArray(self.Status, detachDiskStatus) {
if disk.Status == DISK_INIT {
disk.SetStatus(userCred, DISK_DETACHING, "")
}
taskData := jsonutils.NewDict()
taskData.Add(jsonutils.NewString(diskId), "disk_id")
taskData.Add(jsonutils.NewBool(keepDisk), "keep_disk")
self.GetDriver().StartGuestDetachdiskTask(ctx, userCred, self, taskData, "")
return nil, nil
} else {
return nil, httperrors.NewInvalidStatusError("Server in %s not able to detach disk", self.Status)
}
} else {
return nil, httperrors.NewInvalidStatusError("Disk %s not attached", diskId)
}
}
return nil, httperrors.NewResourceNotFoundError("Disk %s not found", diskId)
}
func (self *SGuest) AllowPerformChangeConfig(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return self.IsOwner(userCred) || self.IsAdmin(userCred)
}
func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
if !utils.IsInStringArray(self.Status, []string{VM_READY}) {
return nil, httperrors.NewInvalidStatusError("Cannot change config in %s", self.Status)
}
if !self.GetDriver().AllowReconfigGuest() {
return nil, httperrors.NewInvalidStatusError("Not allow to change config")
}
host := self.GetHost()
if host == nil {
return nil, httperrors.NewInvalidStatusError("No valid host")
}
var addCpu, addMem int
confs := jsonutils.NewDict()
vcpuCount, err := data.GetString("vcpu_count")
if err == nil {
nVcpu, err := strconv.ParseInt(vcpuCount, 10, 0)
if err != nil {
return nil, httperrors.NewBadRequestError("Params vcpu_count parse error")
}
err = confs.Add(jsonutils.NewInt(nVcpu), "vcpu_count")
if err != nil {
return nil, httperrors.NewBadRequestError("Params vcpu_count parse error")
}
addCpu = int(nVcpu - int64(self.VcpuCount))
}
vmemSize, err := data.GetString("vmem_size")
if err == nil {
if !regutils.MatchSize(vmemSize) {
return nil, httperrors.NewBadRequestError("Memory size must be number[+unit], like 256M, 1G or 256")
}
nVmem, err := fileutils.GetSizeMb(vmemSize, 'M', 1024)
if err != nil {
httperrors.NewBadRequestError("Params vmem_size parse error")
}
err = confs.Add(jsonutils.NewInt(int64(nVmem)), "vmem_size")
if err != nil {
return nil, httperrors.NewBadRequestError("Params vmem_size parse error")
}
addMem = nVmem - self.VmemSize
}
disks := self.GetDisks()
var addDisk int
var diskIdx = 1
var newDiskIdx = 0
var diskSizes = make(map[string]int, 0)
var newDisks = jsonutils.NewDict()
var resizeDisks = jsonutils.NewArray()
for {
diskNum := fmt.Sprintf("disk.%d", diskIdx)
diskDesc, err := data.Get(diskNum)
if err != nil {
break
}
diskConf, err := parseDiskInfo(ctx, userCred, diskDesc)
if err != nil {
return nil, httperrors.NewBadRequestError("Parse disk info error: %s", err)
}
if diskConf.Size > 0 {
if diskIdx >= len(disks) {
newDisks.Add(jsonutils.Marshal(diskConf), fmt.Sprintf("disk.%d", newDiskIdx))
newDiskIdx += 1
addDisk += diskConf.Size
storage := host.GetLeastUsedStorage(diskConf.Backend)
_, ok := diskSizes[storage.Id]
if !ok {
diskSizes[storage.Id] = 0
}
diskSizes[storage.Id] = diskSizes[storage.Id] + diskConf.Size
} else {
disk := disks[diskIdx].GetDisk()
oldSize := disk.DiskSize
if diskConf.Size < oldSize {
return nil, httperrors.NewInputParameterError("Cannot reduce disk size")
} else if diskConf.Size > oldSize {
arr := jsonutils.NewArray(jsonutils.NewString(disks[diskIdx].DiskId), jsonutils.NewInt(int64(diskConf.Size)))
resizeDisks.Add(arr)
addDisk += diskConf.Size - oldSize
storage := disks[diskIdx].GetDisk().GetStorage()
_, ok := diskSizes[storage.Id]
if !ok {
diskSizes[storage.Id] = 0
}
diskSizes[storage.Id] = diskSizes[storage.Id] + diskConf.Size - oldSize
}
}
}
diskIdx += 1
}
for storageId, needSize := range diskSizes {
iStorage, err := StorageManager.FetchById(storageId)
if err != nil {
return nil, httperrors.NewBadRequestError("Fetch storage error: %s", err)
}
storage := iStorage.(*SStorage)
if storage.GetFreeCapacity() < needSize {
return nil, httperrors.NewInsufficientResourceError("Not enough free space")
}
}
if newDisks.Length() > 0 {
confs.Add(newDisks, "create")
}
if resizeDisks.Length() > 0 {
confs.Add(resizeDisks, "resize")
}
if jsonutils.QueryBoolean(data, "auto_start", false) {
confs.Add(jsonutils.NewBool(true), "auto_start")
}
pendingUsage := &SQuota{}
if addCpu > 0 {
pendingUsage.Cpu = addCpu
}
if addMem > 0 {
pendingUsage.Memory = addMem
}
if addDisk > 0 {
pendingUsage.Storage = addDisk
}
if !pendingUsage.IsEmpty() {
err := QuotaManager.CheckSetPendingQuota(ctx, userCred, userCred.GetProjectId(), pendingUsage)
if err != nil {
return nil, httperrors.NewBadRequestError("Check set pending quota error %s", err)
}
}
if newDisks.Length() > 0 {
err := self.CreateDisksOnHost(ctx, userCred, host, newDisks, pendingUsage)
if err != nil {
QuotaManager.CancelPendingUsage(ctx, userCred, self.ProjectId, nil, pendingUsage)
return nil, httperrors.NewBadRequestError("Create disk on host error: %s", err)
}
}
self.StartChangeConfigTask(ctx, userCred, confs, "", pendingUsage)
return nil, nil
}
func (self *SGuest) StartChangeConfigTask(ctx context.Context, userCred mcclient.TokenCredential,
data *jsonutils.JSONDict, parentTaskId string, pendingUsage quotas.IQuota) error {
self.SetStatus(userCred, VM_CHANGE_FLAVOR, "")
task, err := taskman.TaskManager.NewTask(ctx, "GuestChangeConfigTask", self, userCred, data, parentTaskId, "", pendingUsage)
if err != nil {
return err
}
task.ScheduleRun(nil)
return nil
}
func (self *SGuest) DoPendingDelete(ctx context.Context, userCred mcclient.TokenCredential) {
for _, guestdisk := range self.GetDisks() {
disk := guestdisk.GetDisk()
@@ -2039,7 +2287,7 @@ func (self *SGuest) DoPendingDelete(ctx context.Context, userCred mcclient.Token
if utils.IsInStringArray(storage.StorageType, sysutils.LOCAL_STORAGE_TYPES) || disk.DiskType == DISK_TYPE_SYS || disk.DiskType == DISK_TYPE_SWAP {
disk.DoPendingDelete(ctx, userCred)
} else {
self.detachDisk(ctx, disk, userCred)
self.DetachDisk(ctx, disk, userCred)
}
}
self.SVirtualResourceBase.DoPendingDelete(ctx, userCred)
@@ -2059,7 +2307,24 @@ func (self *SGuest) StartUndeployGuestTask(ctx context.Context, userCred mcclien
}
func (self *SGuest) LeaveAllGroups(userCred mcclient.TokenCredential) {
// TODO
groupGuests := make([]SGroupguest, 0)
q := GroupguestManager.Query()
err := q.Filter(sqlchemy.Equals(q.Field("guest_id"), self.Id)).All(&groupGuests)
if err != nil {
log.Errorln(err.Error())
return
}
for _, gg := range groupGuests {
gg.Delete(context.Background(), userCred)
var group SGroup
gq := GroupManager.Query()
err := gq.Filter(sqlchemy.Equals(gq.Field("id"), gg.SrvtagId)).First(&group)
if err != nil {
log.Errorln(err.Error())
return
}
db.OpsLog.LogDetachEvent(self, &group, userCred, nil)
}
}
func (self *SGuest) DetachAllNetworks(ctx context.Context, userCred mcclient.TokenCredential) error {
@@ -2141,6 +2406,29 @@ func (self *SGuest) PerformSyncstatus(ctx context.Context, userCred mcclient.Tok
return nil, err
}
func (self *SGuest) isNotRunningStatus(status string) bool {
if status == VM_READY || status == VM_SUSPEND {
return true
}
return false
}
func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
preStatus := self.Status
_, err := self.SVirtualResourceBase.PerformStatus(ctx, userCred, query, data)
if err != nil {
return nil, err
}
if preStatus != self.Status && !self.isNotRunningStatus(preStatus) && self.isNotRunningStatus(self.Status) {
db.OpsLog.LogEvent(self, db.ACT_STOP, "", userCred)
if self.Status == VM_READY && !self.DisableDelete.Bool() && self.ShutdownBehavior == SHUTDOWN_TERMINATE {
err = self.StartAutoDeleteGuestTask(ctx, userCred, "")
return nil, err
}
}
return nil, nil
}
type SDeployConfig struct {
Path string
Action string
@@ -2561,6 +2849,26 @@ func (self *SGuest) isAllDisksReady() bool {
return ready
}
func (self *SGuest) AllowPerformSuspend(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return self.IsOwner(userCred)
}
func (self *SGuest) PerformSuspend(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
if self.Status == VM_RUNNING {
err := self.StartSuspendTask(ctx, userCred)
return nil, err
}
return nil, httperrors.NewInvalidStatusError("Cannot suspend VM in status %s", self.Status)
}
func (self *SGuest) StartSuspendTask(ctx context.Context, userCred mcclient.TokenCredential) error {
err := self.SetStatus(userCred, VM_SUSPEND, "do suspend")
if err != nil {
return err
}
return self.GetDriver().StartSuspendTask(ctx, userCred, self, nil, "")
}
func (self *SGuest) AllowPerformStart(ctx context.Context,
userCred mcclient.TokenCredential,
query jsonutils.JSONObject,
@@ -2573,14 +2881,13 @@ func (self *SGuest) PerformStart(ctx context.Context, userCred mcclient.TokenCre
if utils.IsInStringArray(self.Status, []string{VM_READY, VM_START_FAILED, VM_SAVE_DISK_FAILED, VM_SUSPEND}) {
if self.isAllDisksReady() {
var kwargs *jsonutils.JSONDict
if data == nil {
if data != nil {
kwargs = data.(*jsonutils.JSONDict)
}
err := self.GetDriver().PerformStart(ctx, userCred, self, kwargs)
return nil, err
} else {
msg := "Some disk not ready"
return nil, httperrors.NewResourceNotReadyError(msg)
return nil, httperrors.NewInvalidStatusError("Some disk not ready")
}
} else {
return nil, httperrors.NewInvalidStatusError("Cannot do start server in status %s", self.Status)
@@ -2645,6 +2952,36 @@ func (self *SGuest) GetDetailsVnc(ctx context.Context, userCred mcclient.TokenCr
}
}
func (self *SGuest) AllowGetDetailsMonitor(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool {
return self.IsOwner(userCred)
}
func (self *SGuest) GetDetailsMonitor(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) {
if utils.IsInStringArray(self.Status, []string{VM_RUNNING, VM_SNAPSHOT_STREAM}) {
cmd, err := query.GetString("command")
if err != nil {
return nil, err
}
return self.SendMonitorCommand(ctx, userCred, cmd)
}
return nil, httperrors.NewInvalidStatusError("Cannot send command in status %s", self.Status)
}
func (self *SGuest) SendMonitorCommand(ctx context.Context, userCred mcclient.TokenCredential, cmd string) (jsonutils.JSONObject, error) {
host := self.GetHost()
url := fmt.Sprintf("%s/servers/%s/monitor", host.ManagerUri, self.Id)
header := http.Header{}
header.Add("X-Auth-Token", userCred.GetTokenString())
body := jsonutils.NewDict()
body.Add(jsonutils.NewString(cmd), "cmd")
_, res, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, body, false)
if err != nil {
return nil, err
}
ret := res.(*jsonutils.JSONDict)
return ret, nil
}
func (model *SGuestManager) AllowPerformCancelDelete(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return userCred.IsSystemAdmin()
}
+27 -1
View File
@@ -2,6 +2,7 @@ package tasks
import (
"context"
"fmt"
"github.com/yunionio/jsonutils"
"github.com/yunionio/log"
@@ -53,6 +54,7 @@ func (self *DiskResizeTask) StartResizeDisk(ctx context.Context, host *models.SH
if err := proc(host, storage, disk, size, self); err != nil {
log.Errorf("request_resize_disk_on_host: %v", err)
self.OnStartResizeDiskFailed(ctx, err)
return
}
self.OnStartResizeDiskSucc(ctx, disk)
}
@@ -69,7 +71,31 @@ func (self *DiskResizeTask) OnStartResizeDiskFailed(ctx context.Context, resion
}
func (self *DiskResizeTask) OnDiskResizeComplete(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) {
disk.SetStatus(self.UserCred, models.DISK_READY, "")
jSize, err := data.Get("disk_size")
if err != nil {
log.Errorf("OnDiskResizeComplete error: %s", err.Error())
self.OnStartResizeDiskFailed(ctx, err)
return
}
size, err := jSize.Int()
if err != nil {
log.Errorf("OnDiskResizeComplete error: %s", err.Error())
self.OnStartResizeDiskFailed(ctx, err)
return
}
oldStatus := disk.Status
_, err = disk.GetModelManager().TableSpec().Update(disk, func() error {
disk.Status = models.DISK_READY
disk.DiskSize = int(size)
return nil
})
if err != nil {
log.Errorf("OnDiskResizeComplete error: %s", err.Error())
self.OnStartResizeDiskFailed(ctx, err)
return
}
notes := fmt.Sprintf("%s=>%s", oldStatus, disk.Status)
db.OpsLog.LogEvent(disk, db.ACT_UPDATE_STATUS, notes, self.UserCred)
self.CleanHostSchedCache(disk)
db.OpsLog.LogEvent(disk, db.ACT_RESIZE, disk.GetShortDesc(), self.UserCred)
self.SetStageComplete(ctx, disk.GetShortDesc())
@@ -0,0 +1,208 @@
package tasks
import (
"context"
"github.com/yunionio/jsonutils"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/lockman"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman"
"github.com/yunionio/onecloud/pkg/compute/models"
)
type GuestChangeConfigTask struct {
SGuestBaseTask
}
func init() {
taskman.RegisterTask(GuestChangeConfigTask{})
}
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.OnDisksResizeComplete(ctx, obj, data)
} else {
guest := obj.(*models.SGuest)
self.DoCreateDisksTask(ctx, guest)
}
}
func (self *GuestChangeConfigTask) OnDisksResizeComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
iResizeDisks, err := self.Params.Get("resize")
if iResizeDisks == nil || err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
resizeDisks := iResizeDisks.(*jsonutils.JSONArray)
for i := 0; i < resizeDisks.Length(); i++ {
iResizeSet, err := resizeDisks.GetAt(i)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
resizeSet := iResizeSet.(*jsonutils.JSONArray)
diskId, err := resizeSet.GetAt(0)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
idStr, err := diskId.GetString()
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
jSize, err := resizeSet.GetAt(1)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
size, err := jSize.Int()
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
iDisk, err := models.DiskManager.FetchById(idStr)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
disk := iDisk.(*models.SDisk)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
if disk.DiskSize < int(size) {
var pendingUsage models.SQuota
err = self.GetPendingUsage(&pendingUsage)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
disk.StartDiskResizeTask(ctx, self.UserCred, size, self.GetTaskId(), &pendingUsage)
return
}
}
guest := obj.(*models.SGuest)
self.DoCreateDisksTask(ctx, guest)
}
func (self *GuestChangeConfigTask) DoCreateDisksTask(ctx context.Context, guest *models.SGuest) {
iCreateData, err := self.Params.Get("create")
if err != nil || iCreateData == nil {
self.OnCreateDisksComplete(ctx, guest, nil)
return
}
data := (iCreateData).(*jsonutils.JSONDict)
self.SetStage("on_create_disks_complete", nil)
guest.StartGuestCreateDiskTask(ctx, self.UserCred, data, self.GetTaskId())
}
func (self *GuestChangeConfigTask) OnCreateDisksComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
iVcpuCount, errCpu := self.Params.Get("vcpu_count")
iVmemSize, errMem := self.Params.Get("vmem_size")
var vcpuCount, vmemSize int64
var err error
guest := obj.(*models.SGuest)
if errCpu == nil || errMem == nil {
if iVcpuCount != nil {
vcpuCount, err = iVcpuCount.Int()
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
}
if iVmemSize != nil {
vmemSize, err = iVmemSize.Int()
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
}
err = guest.GetDriver().RequestChangeVmConfig(ctx, guest, self, vcpuCount, vmemSize)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
var addCpu, addMem = 0, 0
if vcpuCount > 0 {
addCpu = int(vcpuCount - int64(guest.VcpuCount))
if addCpu < 0 {
addCpu = 0
}
}
if vmemSize > 0 {
addMem = int(vmemSize - int64(guest.VmemSize))
if addMem < 0 {
addMem = 0
}
}
_, err = guest.GetModelManager().TableSpec().Update(guest, func() error {
if vcpuCount > 0 {
guest.VcpuCount = int8(vcpuCount)
}
if vmemSize > 0 {
guest.VmemSize = int(vmemSize)
}
return nil
})
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
var pendingUsage models.SQuota
err = self.GetPendingUsage(&pendingUsage)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
// ownerCred := guest.GetOwnerUserCred()
var cancelUsage models.SQuota
if addCpu > 0 {
cancelUsage.Cpu = addCpu
}
if addMem > 0 {
cancelUsage.Memory = addMem
}
lockman.LockClass(ctx, guest.GetModelManager(), guest.ProjectId)
defer lockman.ReleaseClass(ctx, guest.GetModelManager(), guest.ProjectId)
err = models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, guest.ProjectId, &pendingUsage, &cancelUsage)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
err = self.SetPendingUsage(&pendingUsage)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
}
self.SetStage("on_sync_status_complete", nil)
err = guest.StartSyncstatus(ctx, self.UserCred, self.GetTaskId())
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
}
func (self *GuestChangeConfigTask) OnSyncStatusComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
if guest.Status == models.VM_READY && jsonutils.QueryBoolean(self.Params, "auto_start", false) {
self.SetStage("on_guest_start_complete", nil)
guest.StartGueststartTask(ctx, self.UserCred, nil, self.GetTaskId())
} else {
dt := jsonutils.NewDict()
dt.Add(jsonutils.NewString(guest.Id), "id")
self.SetStageComplete(ctx, dt)
}
}
func (self *GuestChangeConfigTask) OnGuestStartComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
dt := jsonutils.NewDict()
dt.Add(jsonutils.NewString(guest.Id), "id")
self.SetStageComplete(ctx, dt)
}
+120
View File
@@ -0,0 +1,120 @@
package tasks
import (
"context"
"fmt"
"github.com/yunionio/jsonutils"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman"
"github.com/yunionio/onecloud/pkg/compute/models"
)
type GuestCreateDiskTask struct {
SGuestBaseTask
}
func (self *GuestCreateDiskTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
self.SetStage("on_disk_prepared", nil)
guest := obj.(*models.SGuest)
err := guest.GetDriver().DoGuestCreateDisksTask(ctx, guest, self)
if err != nil {
self.SetStageFailed(ctx, err.Error())
}
}
func (self *GuestCreateDiskTask) OnDiskPrepared(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
/* --------------------------------------------- */
/* -----------KVMGuestCreateDiskTask------------ */
/* --------------------------------------------- */
type KVMGuestCreateDiskTask struct {
SGuestBaseTask
}
func (self *KVMGuestCreateDiskTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
self.SetStage("on_kvm_disk_prepared", nil)
self.OnKvmDiskPrepared(ctx, obj, data)
}
func (self *KVMGuestCreateDiskTask) OnKvmDiskPrepared(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
var diskIndex = 0
var diskReady = true
for {
diskId, err := self.Params.GetString(fmt.Sprintf("disk.%d.id", diskIndex))
if !diskReady || err != nil {
break
}
iDisk, err := models.DiskManager.FetchById(diskId)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
if iDisk == nil {
self.SetStageFailed(ctx, "Disk not found")
return
}
disk := iDisk.(*models.SDisk)
if disk.Status == models.DISK_INIT {
snapInfo, err := self.Params.GetString(fmt.Sprintf("disk.%d.snapshot", diskIndex))
if err != nil {
snapInfo = ""
}
err = disk.StartDiskCreateTask(ctx, self.UserCred, false, snapInfo, self.GetTaskId())
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
diskReady = false
break
}
diskIndex += 1
}
diskIndex = 0
for {
diskId, err := self.Params.GetString(fmt.Sprintf("disk.%d.id", diskIndex))
if !diskReady || err != nil {
break
}
iDisk, err := models.DiskManager.FetchById(diskId)
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
if iDisk == nil {
self.SetStageFailed(ctx, "Disk not found")
return
}
disk := iDisk.(*models.SDisk)
if disk.Status != models.DISK_READY {
diskReady = false
break
}
diskIndex += 1
}
if diskReady {
guest := obj.(*models.SGuest)
if guest.Status == models.VM_RUNNING {
self.SetStage("on_config_sync_complete", nil)
err := guest.StartSyncstatus(ctx, self.UserCred, self.GetTaskId())
if err != nil {
self.SetStageFailed(ctx, err.Error())
}
} else {
self.SetStageComplete(ctx, nil)
}
}
}
func (self *KVMGuestCreateDiskTask) OnConfigSyncComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
func init() {
taskman.RegisterTask(GuestCreateDiskTask{})
taskman.RegisterTask(KVMGuestCreateDiskTask{})
}
+6 -7
View File
@@ -28,15 +28,14 @@ func (self *GuestDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneModel
func (self *GuestDeleteTask) OnGuestStopComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
guestStatus, _ := self.Params.GetString("guest_status")
if options.Options.EnablePendingDelete && !guest.PendingDeleted &&
!jsonutils.QueryBoolean(self.Params, "purge", false) &&
!jsonutils.QueryBoolean(self.Params, "override_pending_delete", false) {
guestStatus, _ := self.Params.GetString("guest_status")
if !utils.IsInStringArray(guestStatus, []string{models.VM_SCHEDULE_FAILED, models.VM_NETWORK_FAILED, models.VM_DISK_FAILED,
!jsonutils.QueryBoolean(self.Params, "override_pending_delete", false) &&
!utils.IsInStringArray(guestStatus, []string{models.VM_SCHEDULE_FAILED, models.VM_NETWORK_FAILED, models.VM_DISK_FAILED,
models.VM_CREATE_FAILED, models.VM_DEVICE_FAILED}) {
self.StartPendingDeleteGuest(ctx, guest)
return
}
self.StartPendingDeleteGuest(ctx, guest)
return
}
self.OnGuestStopCompleteFailed(ctx, guest, data)
}
@@ -64,6 +63,7 @@ func (self *GuestDeleteTask) OnPendingDeleteComplete(ctx context.Context, obj db
}
func (self *GuestDeleteTask) StartDeleteGuest(ctx context.Context, guest *models.SGuest) {
// No snapshot
self.SetStage("on_guest_detach_disks_complete", nil)
guest.GetDriver().RequestDetachDisksFromGuestForDelete(ctx, guest, self)
}
@@ -100,7 +100,6 @@ func (self *GuestDeleteTask) OnGuestDeleteComplete(ctx context.Context, obj db.I
}
func (self *GuestDeleteTask) DeleteGuest(ctx context.Context, guest *models.SGuest) {
// host := guest.GetHost()
guest.RealDelete(ctx, self.UserCred)
guest.RemoveAllMetadata(ctx, self.UserCred)
db.OpsLog.LogEvent(guest, db.ACT_DELOCATE, nil, self.UserCred)
@@ -4,6 +4,7 @@ import (
"context"
"github.com/yunionio/jsonutils"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman"
"github.com/yunionio/onecloud/pkg/compute/models"
@@ -2,10 +2,14 @@ package tasks
import (
"context"
"fmt"
"github.com/yunionio/jsonutils"
"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman"
"github.com/yunionio/onecloud/pkg/compute/models"
"github.com/yunionio/pkg/utils"
)
type GuestDetachDiskTask struct {
@@ -17,5 +21,83 @@ func init() {
}
func (self *GuestDetachDiskTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
diskId, _ := self.Params.GetString("disk_id")
objDisk, err := models.DiskManager.FetchById(diskId)
if err != nil {
self.OnTaskFail(ctx, guest, err)
return
}
disk := objDisk.(*models.SDisk)
if disk == nil {
self.OnTaskFail(ctx, guest, fmt.Errorf("Connot find disk %s", diskId))
return
}
guest.DetachDisk(ctx, disk, self.UserCred)
if disk.Status == models.DISK_INIT {
self.OnSyncConfigComplete(ctx, guest, nil)
return
}
host := guest.GetHost()
purge := false
if host != nil && host.Status == models.HOST_DISABLED && jsonutils.QueryBoolean(self.Params, "purge", false) {
purge = true
}
detachStatus, err := guest.GetDriver().GetDetachDiskStatus()
if err != nil {
self.OnTaskFail(ctx, guest, 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)
}
}
func (self *GuestDetachDiskTask) OnSyncConfigComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) {
diskId, _ := self.Params.GetString("disk_id")
objDisk, err := models.DiskManager.FetchById(diskId)
if err != nil {
self.OnTaskFail(ctx, guest, err)
return
}
disk := objDisk.(*models.SDisk)
if disk == nil {
self.OnTaskFail(ctx, guest, fmt.Errorf("Connot find disk %s", diskId))
return
}
keepDisk := jsonutils.QueryBoolean(self.Params, "keep_disk", true)
host := guest.GetHost()
purge := false
if host != nil && host.Status == models.HOST_DISABLED && jsonutils.QueryBoolean(self.Params, "purge", false) {
purge = true
}
if disk.Status == models.DISK_INIT {
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 {
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)
return
}
} else {
self.SetStageComplete(ctx, nil)
}
}
func (self *GuestDetachDiskTask) OnTaskFail(ctx context.Context, guest *models.SGuest, err error) {
self.SetStageFailed(ctx, err.Error())
log.Errorf("Guest %s GuestDetachDiskTask failed %s", guest.Id, err.Error())
}
func (self *GuestDetachDiskTask) OnDiskDeleteComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
self.SetStageComplete(ctx, nil)
}
+1
View File
@@ -6,6 +6,7 @@ import (
"github.com/yunionio/jsonutils"
"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman"
"github.com/yunionio/onecloud/pkg/compute/models"
+47
View File
@@ -0,0 +1,47 @@
package tasks
import (
"context"
"github.com/yunionio/jsonutils"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman"
"github.com/yunionio/onecloud/pkg/compute/models"
)
type GuestSuspendTask struct {
SGuestBaseTask
}
func init() {
taskman.RegisterTask(GuestSuspendTask{})
}
func (self *GuestSuspendTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
db.OpsLog.LogEvent(guest, db.ACT_STOPPING, "", self.UserCred)
guest.SetStatus(self.UserCred, models.VM_SUSPENDING, "GuestSusPendTask")
self.SetStage("on_suspend_complete", nil)
err := guest.GetDriver().RqeuestSuspendOnHost(ctx, guest, self)
if err != nil {
self.OnSuspendGuestFail(guest, err.Error())
}
}
func (self *GuestSuspendTask) OnSuspendComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
guest.SetStatus(self.UserCred, models.VM_SUSPEND, "")
db.OpsLog.LogEvent(guest, db.ACT_STOP, "", self.UserCred)
self.SetStageComplete(ctx, nil)
}
func (self *GuestSuspendTask) OnSuspendCompleteFailed(ctx context.Context, obj db.IStandaloneModel, err jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
guest.SetStatus(self.UserCred, models.VM_RUNNING, "")
db.OpsLog.LogEvent(guest, db.ACT_STOP_FAIL, err.String(), self.UserCred)
}
func (self *GuestSuspendTask) OnSuspendGuestFail(guest *models.SGuest, reason string) {
guest.SetStatus(self.UserCred, models.VM_SUSPEND_FAILED, reason)
}
+3 -1
View File
@@ -48,7 +48,9 @@ func (self *GuestSyncstatusTask) OnGetStatusSucc(ctx context.Context, guest *mod
default:
statusStr = models.VM_UNKNOWN
}
guest.SetStatus(self.UserCred, statusStr, "syncstatus")
statusData := jsonutils.NewDict()
statusData.Add(jsonutils.NewString(statusStr), "status")
guest.PerformStatus(ctx, self.UserCred, nil, statusData)
self.SetStageComplete(ctx, nil)
}
+1 -2
View File
@@ -4,6 +4,7 @@ import (
"context"
"github.com/yunionio/jsonutils"
"github.com/yunionio/onecloud/pkg/cloudcommon/db"
"github.com/yunionio/onecloud/pkg/cloudcommon/db/taskman"
"github.com/yunionio/onecloud/pkg/compute/models"
@@ -33,8 +34,6 @@ func (self *GuestUndeployTask) OnInit(ctx context.Context, obj db.IStandaloneMod
err := guest.GetDriver().RequestUndeployGuestOnHost(ctx, guest, host, self)
if err != nil {
self.OnStartDeleteGuestFail(ctx, err)
} else {
// do nothing
}
} else {
self.SetStageComplete(ctx, nil)