From f04463b5269157e03739412952ee0c38c2688bac Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Wed, 1 Aug 2018 22:18:32 +0800 Subject: [PATCH] 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 --- cmd/climc/shell/servers.go | 2 +- pkg/appsrv/dispatcher/dispatcher.go | 10 +- pkg/cloudcommon/db/db_dispatcher.go | 4 +- pkg/cloudcommon/db/quotas/quotas.go | 4 +- pkg/cloudcommon/db/taskman/tasks.go | 8 +- pkg/cloudcommon/db/virtualresource.go | 4 +- pkg/compute/guestdrivers/baremetals.go | 15 +- pkg/compute/guestdrivers/base.go | 35 +- pkg/compute/guestdrivers/esxi.go | 17 + pkg/compute/guestdrivers/kvm.go | 90 ++++- pkg/compute/guestdrivers/virtualization.go | 23 ++ pkg/compute/models/disks.go | 15 +- pkg/compute/models/guestdrivers.go | 16 +- pkg/compute/models/guestnetworks.go | 2 + pkg/compute/models/guests.go | 355 +++++++++++++++++- pkg/compute/tasks/disk_resize_task.go | 28 +- pkg/compute/tasks/guest_change_config_task.go | 208 ++++++++++ pkg/compute/tasks/guest_create_disk_task.go | 120 ++++++ pkg/compute/tasks/guest_delete_task.go | 13 +- .../tasks/guest_detach_all_disks_task.go | 1 + pkg/compute/tasks/guest_detach_disk_task.go | 82 ++++ pkg/compute/tasks/guest_stop_task.go | 1 + pkg/compute/tasks/guest_suspend_task.go | 47 +++ pkg/compute/tasks/guest_syncstatus_task.go | 4 +- pkg/compute/tasks/guest_undeploy_task.go | 3 +- 25 files changed, 1048 insertions(+), 59 deletions(-) create mode 100644 pkg/compute/tasks/guest_change_config_task.go create mode 100644 pkg/compute/tasks/guest_create_disk_task.go create mode 100644 pkg/compute/tasks/guest_suspend_task.go diff --git a/cmd/climc/shell/servers.go b/cmd/climc/shell/servers.go index 9a5fc8df52..038139983c 100644 --- a/cmd/climc/shell/servers.go +++ b/cmd/climc/shell/servers.go @@ -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 diff --git a/pkg/appsrv/dispatcher/dispatcher.go b/pkg/appsrv/dispatcher/dispatcher.go index c6356128d2..6bfed3aee8 100644 --- a/pkg/appsrv/dispatcher/dispatcher.go +++ b/pkg/appsrv/dispatcher/dispatcher.go @@ -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() } diff --git a/pkg/cloudcommon/db/db_dispatcher.go b/pkg/cloudcommon/db/db_dispatcher.go index ad5e5cbe17..c3ce3c6b14 100644 --- a/pkg/cloudcommon/db/db_dispatcher.go +++ b/pkg/cloudcommon/db/db_dispatcher.go @@ -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) diff --git a/pkg/cloudcommon/db/quotas/quotas.go b/pkg/cloudcommon/db/quotas/quotas.go index d09b6c1007..da6bae02cb 100644 --- a/pkg/cloudcommon/db/quotas/quotas.go +++ b/pkg/cloudcommon/db/quotas/quotas.go @@ -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 { diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index 0ebf6cb7eb..e5bfead6cf 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -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) diff --git a/pkg/cloudcommon/db/virtualresource.go b/pkg/cloudcommon/db/virtualresource.go index a07ac9e192..7ee602e30c 100644 --- a/pkg/cloudcommon/db/virtualresource.go +++ b/pkg/cloudcommon/db/virtualresource.go @@ -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 { diff --git a/pkg/compute/guestdrivers/baremetals.go b/pkg/compute/guestdrivers/baremetals.go index 82828aad61..09751c0944 100644 --- a/pkg/compute/guestdrivers/baremetals.go +++ b/pkg/compute/guestdrivers/baremetals.go @@ -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") +} diff --git a/pkg/compute/guestdrivers/base.go b/pkg/compute/guestdrivers/base.go index ca374ec2a0..948bbd0c38 100644 --- a/pkg/compute/guestdrivers/base.go +++ b/pkg/compute/guestdrivers/base.go @@ -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") +} diff --git a/pkg/compute/guestdrivers/esxi.go b/pkg/compute/guestdrivers/esxi.go index 6d2536025d..e516d0c266 100644 --- a/pkg/compute/guestdrivers/esxi.go +++ b/pkg/compute/guestdrivers/esxi.go @@ -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 +} diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index b1fe945745..73a6439786 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -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 +} diff --git a/pkg/compute/guestdrivers/virtualization.go b/pkg/compute/guestdrivers/virtualization.go index fc0b1a8a8f..11d62fe15a 100644 --- a/pkg/compute/guestdrivers/virtualization.go +++ b/pkg/compute/guestdrivers/virtualization.go @@ -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 +} diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 67cdebd1b7..5b50f79363 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -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 } diff --git a/pkg/compute/models/guestdrivers.go b/pkg/compute/models/guestdrivers.go index dd5feb0b9d..3370ef5529 100644 --- a/pkg/compute/models/guestdrivers.go +++ b/pkg/compute/models/guestdrivers.go @@ -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 diff --git a/pkg/compute/models/guestnetworks.go b/pkg/compute/models/guestnetworks.go index 002547a40c..f401930339 100644 --- a/pkg/compute/models/guestnetworks.go +++ b/pkg/compute/models/guestnetworks.go @@ -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 { diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 584026baa9..7a228421b8 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -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() } diff --git a/pkg/compute/tasks/disk_resize_task.go b/pkg/compute/tasks/disk_resize_task.go index eb6b98af37..fb20ca0fcb 100644 --- a/pkg/compute/tasks/disk_resize_task.go +++ b/pkg/compute/tasks/disk_resize_task.go @@ -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()) diff --git a/pkg/compute/tasks/guest_change_config_task.go b/pkg/compute/tasks/guest_change_config_task.go new file mode 100644 index 0000000000..e75b37bd06 --- /dev/null +++ b/pkg/compute/tasks/guest_change_config_task.go @@ -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) +} diff --git a/pkg/compute/tasks/guest_create_disk_task.go b/pkg/compute/tasks/guest_create_disk_task.go new file mode 100644 index 0000000000..4ba43f712a --- /dev/null +++ b/pkg/compute/tasks/guest_create_disk_task.go @@ -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{}) +} diff --git a/pkg/compute/tasks/guest_delete_task.go b/pkg/compute/tasks/guest_delete_task.go index 6c18a50058..2139a92d94 100644 --- a/pkg/compute/tasks/guest_delete_task.go +++ b/pkg/compute/tasks/guest_delete_task.go @@ -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) diff --git a/pkg/compute/tasks/guest_detach_all_disks_task.go b/pkg/compute/tasks/guest_detach_all_disks_task.go index cfc6465322..a25d887fb5 100644 --- a/pkg/compute/tasks/guest_detach_all_disks_task.go +++ b/pkg/compute/tasks/guest_detach_all_disks_task.go @@ -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" diff --git a/pkg/compute/tasks/guest_detach_disk_task.go b/pkg/compute/tasks/guest_detach_disk_task.go index 2dff6ca781..e07484a62f 100644 --- a/pkg/compute/tasks/guest_detach_disk_task.go +++ b/pkg/compute/tasks/guest_detach_disk_task.go @@ -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) } diff --git a/pkg/compute/tasks/guest_stop_task.go b/pkg/compute/tasks/guest_stop_task.go index b2843cc377..e5dc2d4de4 100644 --- a/pkg/compute/tasks/guest_stop_task.go +++ b/pkg/compute/tasks/guest_stop_task.go @@ -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" diff --git a/pkg/compute/tasks/guest_suspend_task.go b/pkg/compute/tasks/guest_suspend_task.go new file mode 100644 index 0000000000..348d6b24b6 --- /dev/null +++ b/pkg/compute/tasks/guest_suspend_task.go @@ -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) +} diff --git a/pkg/compute/tasks/guest_syncstatus_task.go b/pkg/compute/tasks/guest_syncstatus_task.go index 201c20cd66..7ae86fdfe0 100644 --- a/pkg/compute/tasks/guest_syncstatus_task.go +++ b/pkg/compute/tasks/guest_syncstatus_task.go @@ -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) } diff --git a/pkg/compute/tasks/guest_undeploy_task.go b/pkg/compute/tasks/guest_undeploy_task.go index 90958eb21d..d1e845d384 100644 --- a/pkg/compute/tasks/guest_undeploy_task.go +++ b/pkg/compute/tasks/guest_undeploy_task.go @@ -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)