Merge pull request #12049 from ioito/automated-cherry-pick-of-#11987-upstream-release-3.8

Automated cherry pick of #11987: feat(region): cloudpods operation
This commit is contained in:
Zexi Li
2021-08-27 16:49:49 +08:00
committed by GitHub
36 changed files with 689 additions and 550 deletions
-13
View File
@@ -543,19 +543,6 @@ type ServerCloneInput struct {
PreferHostId string `json:"prefer_host_id"`
}
type ServerDeployInput struct {
apis.Meta
Id string
Keypair string `json:"keypair"`
DeleteKeypair *bool `json:"__delete_keypair__"`
DeployConfigs []*DeployConfig `json:"deploy_configs"`
ResetPassword *bool `json:"reset_password"`
Password string `json:"password"`
AutoStart *bool `json:"auto_start"`
}
type GuestBatchMigrateRequest struct {
apis.Meta
+87 -24
View File
@@ -121,12 +121,14 @@ func (input *ServerListInput) AfterUnmarshal() {
type ServerRebuildRootInput struct {
apis.Meta
// 镜像名称
Image string `json:"image"`
// swagger: ignore
Image string `json:"image" yunion-deprecated-by:"image_id"`
// 镜像 id
// required: true
ImageId string `json:"image_id"`
Keypair string `json:"keypair"`
ImageId string `json:"image_id"`
// swagger: ignore
Keypair string `json:"keypair" yunion-deprecated-by:"keypair_id"`
// 秘钥Id
KeypairId string `json:"keypair_id"`
ResetPassword *bool `json:"reset_password"`
Password string `json:"password"`
@@ -134,26 +136,6 @@ type ServerRebuildRootInput struct {
AllDisks *bool `json:"all_disks"`
}
func (i ServerRebuildRootInput) GetImageName() string {
if len(i.Image) > 0 {
return i.Image
}
if len(i.ImageId) > 0 {
return i.ImageId
}
return ""
}
func (i ServerRebuildRootInput) GetKeypairName() string {
if len(i.Keypair) > 0 {
return i.Keypair
}
if len(i.KeypairId) > 0 {
return i.KeypairId
}
return ""
}
type ServerResumeInput struct {
apis.Meta
}
@@ -555,3 +537,84 @@ type ServerMigrateNetworkInput struct {
// Destination network Id
Dest string `json:"dest"`
}
type ServerDeployInput struct {
apis.Meta
// swagger: ignore
Keypair string `json:"keypair" yunion-deprecated-by:"keypair_id"`
// 秘钥Id
KeypairId string `json:"keypair_id"`
// 清理指定公钥
// 若指定的秘钥Id和虚拟机的秘钥Id不相同, 则清理旧的公钥
DeletePublicKey string `json:"delete_public_key"`
// 解绑当前虚拟机秘钥, 并清理公钥信息
DeleteKeypair bool `json:"__delete_keypair__"`
// 生成随机密码, 优先级低于password
ResetPassword bool `json:"reset_password"`
// 重置指定密码
Password string `json:"password"`
// 部署完成后是否自动启动
// 若虚拟机重置密码后需要重启生效,并且当前虚拟机状态为running, 此参数默认为true
// 若虚拟机状态为ready, 指定此参数后,部署完成后,虚拟机会自动启动
AutoStart bool `json:"auto_start"`
// swagger: ignore
Restart bool `json:"restart"`
// swagger: ignore
DeployConfigs []*DeployConfig `json:"deploy_configs"`
}
type ServerUserDataInput struct {
UserData string `json:"user_data"`
}
type ServerAttachDiskInput struct {
DiskId string `json:"disk_id"`
}
type ServerDetachDiskInput struct {
// 磁盘Id,若磁盘未挂载在虚拟机上,不返回错误
DiskId string `json:"disk_id"`
// 是否保留磁盘
// default: false
KeepDisk bool `json:"keep_disk"`
}
type ServerChangeConfigInput struct {
// 实例类型, 优先级高于vcpu_count和vmem_size
InstanceType string `json:"instance_type"`
// swagger: ignore
Sku string `json:"sku" yunion-deprecated-by:"instance_type"`
// swagger: ignore
Flavor string `json:"flavor" yunion-deprecated-by:"instance_type"`
// cpu大小
VcpuCount int `json:"vcpu_count"`
// 内存大小, 1024M, 1G
VmemSize string `json:"vmem_size"`
// 调整完配置后是否自动启动
AutoStart bool `json:"auto_start"`
Disks []DiskConfig `json:"disks"`
}
type ServerUpdateInput struct {
apis.VirtualResourceBaseUpdateInput
// 删除保护开关
DisableDelete *bool `json:"disable_delete"`
// 启动顺序
BootOrder *string `json:"boot_order"`
// 关机执行操作
ShutdownBehavior *string `json:"shutdown_behavior"`
Vga *string `json:"vga"`
Vdi *string `json:"vdi"`
Machine *string `json:"machine"`
Bios *string `json:"bios"`
SrcIpCheck *bool `json:"src_ip_check"`
SrcMacCheck *bool `json:"src_mac_check"`
}
+1
View File
@@ -32,4 +32,5 @@ const (
HUAWEI = "huawei"
APSARA = "apsara"
JDCLOUD = "jdcloud"
CLOUDPODS = "cloudpods"
)
+7
View File
@@ -20,3 +20,10 @@ type SNetworkCreateOptions struct {
ProjectId string
Cidr string
}
type SWireCreateOptions struct {
Name string
ZoneId string
Bandwidth int
Mtu int
}
+1 -1
View File
@@ -552,8 +552,8 @@ type ICloudVpc interface {
GetRegion() ICloudRegion
GetIsDefault() bool
GetCidrBlock() string
// GetStatus() string
GetIWires() ([]ICloudWire, error)
CreateIWire(opts *SWireCreateOptions) (ICloudWire, error)
GetISecurityGroups() ([]ICloudSecurityGroup, error)
GetIRouteTables() ([]ICloudRouteTable, error)
GetIRouteTableById(routeTableId string) (ICloudRouteTable, error)
-8
View File
@@ -19,7 +19,6 @@ import (
"fmt"
"strings"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
@@ -207,13 +206,6 @@ func (self *SAzureGuestDriver) ValidateCreateData(ctx context.Context, userCred
return input, nil
}
func (self *SAzureGuestDriver) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
if data.Contains("name") {
return nil, httperrors.NewInputParameterError("cannot support change azure instance name")
}
return data, nil
}
func (self *SAzureGuestDriver) GetGuestInitialStateAfterCreate() string {
return api.VM_RUNNING
}
+4 -4
View File
@@ -323,13 +323,14 @@ func (self *SBaremetalGuestDriver) RequestGuestCreateInsertIso(ctx context.Conte
return guest.StartInsertIsoTask(ctx, imageId, true, guest.HostId, task.GetUserCred(), task.GetTaskId())
}
func (self *SBaremetalGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) (jsonutils.JSONObject, error) {
func (self *SBaremetalGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error {
desc := guest.GetJsonDescAtBaremetal(ctx, host)
config := jsonutils.NewDict()
config.Set("desc", desc)
headers := task.GetTaskRequestHeader()
url := fmt.Sprintf("/baremetals/%s/servers/%s/start", host.Id, guest.Id)
return host.BaremetalSyncRequest(ctx, "POST", url, headers, config)
_, err := host.BaremetalSyncRequest(ctx, "POST", url, headers, config)
return err
}
func (self *SBaremetalGuestDriver) RequestStopGuestForDelete(ctx context.Context, guest *models.SGuest,
@@ -376,8 +377,7 @@ func (self *SBaremetalGuestDriver) StartGuestStopTask(guest *models.SGuest, ctx
if err != nil {
return err
}
task.ScheduleRun(nil)
return nil
return task.ScheduleRun(nil)
}
func (self *SBaremetalGuestDriver) RequestUndeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
-4
View File
@@ -178,10 +178,6 @@ func (self *SBaseGuestDriver) ValidateResizeDisk(guest *models.SGuest, disk *mod
return fmt.Errorf("This Guest driver dose not implement ValidateResizeDisk")
}
func (self *SBaseGuestDriver) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
return data, nil
}
func (self *SBaseGuestDriver) GetDeployStatus() ([]string, error) {
return []string{}, fmt.Errorf("This Guest driver dose not implement GetDeployStatus")
}
+36
View File
@@ -49,6 +49,42 @@ func (self *SCloudpodsGuestDriver) GetProvider() string {
return api.CLOUD_PROVIDER_CLOUDPODS
}
func (self *SCloudpodsGuestDriver) GetGuestInitialStateAfterCreate() string {
return api.VM_READY
}
func (self *SCloudpodsGuestDriver) GetDetachDiskStatus() ([]string, error) {
return []string{api.VM_READY, api.VM_RUNNING}, nil
}
func (self *SCloudpodsGuestDriver) GetAttachDiskStatus() ([]string, error) {
return []string{api.VM_READY, api.VM_RUNNING}, nil
}
func (self *SCloudpodsGuestDriver) GetRebuildRootStatus() ([]string, error) {
return []string{api.VM_READY, api.VM_RUNNING}, nil
}
func (self *SCloudpodsGuestDriver) GetChangeConfigStatus(guest *models.SGuest) ([]string, error) {
return []string{api.VM_READY, api.VM_RUNNING}, nil
}
func (self *SCloudpodsGuestDriver) GetDeployStatus() ([]string, error) {
return []string{api.VM_READY, api.VM_ADMIN}, nil
}
func (self *SCloudpodsGuestDriver) IsSupportCdrom(guest *models.SGuest) (bool, error) {
return true, nil
}
func (self *SCloudpodsGuestDriver) IsSupportMigrate() bool {
return true
}
func (self *SCloudpodsGuestDriver) IsSupportLiveMigrate() bool {
return true
}
func (self *SCloudpodsGuestDriver) GetComputeQuotaKeys(scope rbacutils.TRbacScope, ownerId mcclient.IIdentityProvider, brand string) models.SComputeResourceKeys {
keys := models.SComputeResourceKeys{}
keys.SBaseProjectQuotaKeys = quotas.OwnerIdProjectQuotaKeys(scope, ownerId)
+2 -2
View File
@@ -96,8 +96,8 @@ func (self *SContainerDriver) RequestGuestHotAddIso(ctx context.Context, guest *
return nil
}
func (self *SContainerDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) (jsonutils.JSONObject, error) {
return nil, httperrors.NewUnsupportOperationError("")
func (self *SContainerDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error {
return httperrors.NewUnsupportOperationError("")
}
func (self *SContainerDriver) RequestStopOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
+6 -14
View File
@@ -199,22 +199,17 @@ func (self *SGoogleGuestDriver) GetUserDataType() string {
return cloudprovider.CLOUD_SHELL_WITHOUT_ENCRYPT
}
func (self *SGoogleGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) (jsonutils.JSONObject, error) {
ihost, err := host.GetIHost()
func (self *SGoogleGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error {
ivm, err := guest.GetIVM()
if err != nil {
return nil, errors.Wrap(err, "host.GetIHost")
}
ivm, err := ihost.GetIVMById(guest.GetExternalId())
if err != nil {
return nil, errors.Wrap(err, "ihost.GetIVMById")
return errors.Wrap(err, "GetIVM")
}
result := jsonutils.NewDict()
if ivm.GetStatus() != api.VM_RUNNING {
err := ivm.StartVM(ctx)
if err != nil {
return nil, errors.Wrap(err, "ivm.StartVM")
return errors.Wrap(err, "ivm.StartVM")
}
vm := ivm.(*google.SInstance)
updateUserdata := false
@@ -247,12 +242,9 @@ func (self *SGoogleGuestDriver) RequestStartOnHost(ctx context.Context, guest *m
log.Errorf("failed to update google userdata")
}
}
task.ScheduleRun(result)
} else {
result.Add(jsonutils.NewBool(true), "is_running")
return task.ScheduleRun(result)
}
return result, nil
return guest.SetStatus(userCred, api.VM_RUNNING, "StartOnHost")
}
func (self *SGoogleGuestDriver) RemoteActionAfterGuestCreated(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, host *models.SHost, iVM cloudprovider.ICloudVM, desc *cloudprovider.SManagedVMCreateConfig) {
+5 -5
View File
@@ -283,13 +283,13 @@ func (self *SKVMGuestDriver) OnGuestDeployTaskDataReceived(ctx context.Context,
return nil
}
func (self *SKVMGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) (jsonutils.JSONObject, error) {
func (self *SKVMGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error {
header := self.getTaskRequestHeader(task)
config := jsonutils.NewDict()
desc, err := guest.GetDriver().GetJsonDescAtHost(ctx, userCred, guest, host, nil)
if err != nil {
return nil, errors.Wrapf(err, "GetJsonDescAtHost")
return errors.Wrapf(err, "GetJsonDescAtHost")
}
config.Add(desc, "desc")
params := task.GetParams()
@@ -297,11 +297,11 @@ func (self *SKVMGuestDriver) RequestStartOnHost(ctx context.Context, guest *mode
config.Add(params, "params")
}
url := fmt.Sprintf("%s/servers/%s/start", host.ManagerUri, guest.Id)
_, res, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, config, false)
_, _, err = httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, config, false)
if err != nil {
return nil, err
return err
}
return res, nil
return nil
}
func (self *SKVMGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) {
+17 -17
View File
@@ -315,29 +315,27 @@ func (self *SManagedVirtualizedGuestDriver) RequestAttachDisk(ctx context.Contex
return nil
}
func (self *SManagedVirtualizedGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) (jsonutils.JSONObject, error) {
ihost, err := host.GetIHost()
func (self *SManagedVirtualizedGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error {
ivm, err := guest.GetIVM()
if err != nil {
return nil, err
}
ivm, err := ihost.GetIVMById(guest.GetExternalId())
if err != nil {
return nil, err
return errors.Wrapf(err, "GetIVM")
}
result := jsonutils.NewDict()
if ivm.GetStatus() != api.VM_RUNNING {
err := ivm.StartVM(ctx)
if err != nil {
return nil, err
return errors.Wrapf(err, "StartVM")
}
task.ScheduleRun(result)
} else {
result.Add(jsonutils.NewBool(true), "is_running")
err = cloudprovider.WaitStatus(ivm, api.VM_RUNNING, time.Second*5, time.Minute*10)
if err != nil {
return errors.Wrapf(err, "Wait vm running")
}
// 虚拟机开机,公网ip自动生成
guest.SyncAllWithCloudVM(ctx, userCred, host, ivm)
return task.ScheduleRun(result)
}
return result, nil
return guest.SetStatus(userCred, api.VM_RUNNING, "StartOnHost")
}
func (self *SManagedVirtualizedGuestDriver) RequestDeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error {
@@ -614,9 +612,8 @@ func (self *SManagedVirtualizedGuestDriver) RemoteDeployGuestForDeploy(ctx conte
func (self *SManagedVirtualizedGuestDriver) RemoteDeployGuestForRebuildRoot(ctx context.Context, guest *models.SGuest, ihost cloudprovider.ICloudHost, task taskman.ITask, desc cloudprovider.SManagedVMCreateConfig) (jsonutils.JSONObject, error) {
iVM, err := ihost.GetIVMById(guest.GetExternalId())
if err != nil || iVM == nil {
log.Errorf("cannot find vm %s", err)
return nil, fmt.Errorf("cannot find vm")
if err != nil {
return nil, errors.Wrapf(err, "ihost.GetIVMById(%s)", guest.GetExternalId())
}
if len(desc.UserData) > 0 {
@@ -624,6 +621,7 @@ func (self *SManagedVirtualizedGuestDriver) RemoteDeployGuestForRebuildRoot(ctx
if err != nil {
log.Errorf("update userdata fail %s", err)
}
cloudprovider.WaitMultiStatus(iVM, []string{api.VM_READY, api.VM_RUNNING}, time.Second*5, time.Minute*3)
}
diskId, err := func() (string, error) {
@@ -757,6 +755,8 @@ func (self *SManagedVirtualizedGuestDriver) RequestStopOnHost(ctx context.Contex
if err != nil {
return nil, errors.Wrapf(err, "wait server stop after 5 miniutes")
}
// 公有云关机,公网ip会释放
guest.SyncAllWithCloudVM(ctx, task.GetUserCred(), host, ivm)
return nil, nil
})
return nil
+1 -2
View File
@@ -225,8 +225,7 @@ func (self *SVirtualizedGuestDriver) StartGuestStopTask(guest *models.SGuest, ct
if err != nil {
return err
}
task.ScheduleRun(nil)
return nil
return task.ScheduleRun(nil)
}
func (self *SVirtualizedGuestDriver) StartGuestResetTask(guest *models.SGuest, ctx context.Context, userCred mcclient.TokenCredential, isHard bool, parentTaskId string) error {
+137 -219
View File
@@ -182,17 +182,17 @@ func (self *SGuest) AllowPerformSaveImage(ctx context.Context, userCred mcclient
return self.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, self, "save-image")
}
func (self *SGuest) PerformSaveImage(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerSaveImageInput) (jsonutils.JSONObject, error) {
func (self *SGuest) PerformSaveImage(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerSaveImageInput) (api.ServerSaveImageInput, error) {
if !utils.IsInStringArray(self.Status, []string{api.VM_READY, api.VM_RUNNING}) {
return nil, httperrors.NewInputParameterError("Cannot save image in status %s", self.Status)
return input, httperrors.NewInputParameterError("Cannot save image in status %s", self.Status)
}
input.Restart = (self.Status == api.VM_RUNNING) || input.AutoStart
if len(input.Name) == 0 && len(input.GenerateName) == 0 {
return nil, httperrors.NewInputParameterError("Image name is required")
return input, httperrors.NewInputParameterError("Image name is required")
}
disks := self.CategorizeDisks()
if disks.Root == nil {
return nil, httperrors.NewInputParameterError("No root image")
return input, httperrors.NewInputParameterError("No root image")
}
input.OsType = self.OsType
if len(input.OsType) == 0 {
@@ -214,14 +214,14 @@ func (self *SGuest) PerformSaveImage(ctx context.Context, userCred mcclient.Toke
var err error
input.ImageId, err = disks.Root.PrepareSaveImage(ctx, userCred, input)
if err != nil {
return nil, errors.Wrapf(err, "PrepareSaveImage")
return input, errors.Wrapf(err, "PrepareSaveImage")
}
}
if len(input.Name) == 0 {
input.Name = input.GenerateName
}
return nil, self.StartGuestSaveImage(ctx, userCred, input, "")
return input, self.StartGuestSaveImage(ctx, userCred, input, "")
}
func (self *SGuest) StartGuestSaveImage(ctx context.Context, userCred mcclient.TokenCredential, input api.ServerSaveImageInput, parentTaskId string) error {
@@ -653,71 +653,47 @@ func (self *SGuest) GetOldPassword(ctx context.Context, userCred mcclient.TokenC
return password
}
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")
}
func (self *SGuest) PerformDeploy(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerDeployInput) (jsonutils.JSONObject, error) {
self.saveOldPassword(ctx, userCred)
if kwargs.Contains("__delete_keypair__") || kwargs.Contains("keypair") {
var kpId string
if kwargs.Contains("keypair") {
keypair, _ := kwargs.GetString("keypair")
iKp, err := KeypairManager.FetchByIdOrName(userCred, keypair)
if input.DeleteKeypair || len(input.KeypairId) > 0 {
if len(input.KeypairId) > 0 {
_, err := validators.ValidateModel(userCred, KeypairManager, &input.KeypairId)
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 {
if self.KeypairId != input.KeypairId {
okey := self.getKeypair()
if okey != nil {
kwargs.Set("delete_public_key", jsonutils.NewString(okey.PublicKey))
input.DeletePublicKey = okey.PublicKey
}
diff, err := db.Update(self, func() error {
self.KeypairId = kpId
self.KeypairId = input.KeypairId
return nil
})
if err != nil {
log.Errorf("update keypair fail: %s", err)
return nil, httperrors.NewInternalServerError("%v", err)
return nil, httperrors.NewInternalServerError("update keypairId %v", err)
}
db.OpsLog.LogEvent(self, db.ACT_UPDATE, diff, userCred)
kwargs.Set("reset_password", jsonutils.JSONTrue)
input.ResetPassword = true
}
}
var resetPasswd bool
passwdStr, _ := kwargs.GetString("password")
if len(passwdStr) > 0 {
err := seclib2.ValidatePassword(passwdStr)
if len(input.Password) > 0 {
err := seclib2.ValidatePassword(input.Password)
if err != nil {
return nil, err
}
resetPasswd = true
} else {
resetPasswd = jsonutils.QueryBoolean(kwargs, "reset_password", false)
}
if resetPasswd {
kwargs.Set("reset_password", jsonutils.JSONTrue)
} else {
kwargs.Set("reset_password", jsonutils.JSONFalse)
input.ResetPassword = true
}
// 变更密码/密钥时需要Restart才能生效。更新普通字段不需要Restart, Azure需要在运行状态下操作
doRestart := false
if resetPasswd {
if input.ResetPassword {
doRestart = self.GetDriver().IsNeedRestartForResetLoginInfo()
}
@@ -727,13 +703,13 @@ func (self *SGuest) PerformDeploy(ctx context.Context, userCred mcclient.TokenCr
}
if utils.IsInStringArray(self.Status, deployStatus) {
if (doRestart && self.Status == api.VM_RUNNING) || (self.Status != api.VM_RUNNING && (jsonutils.QueryBoolean(kwargs, "auto_start", false) || jsonutils.QueryBoolean(kwargs, "restart", false))) {
kwargs.Set("restart", jsonutils.JSONTrue)
if (doRestart && self.Status == api.VM_RUNNING) || (self.Status != api.VM_RUNNING && (input.AutoStart || input.Restart)) {
input.Restart = true
} else {
// 避免前端直接传restart参数, 越过校验
kwargs.Set("restart", jsonutils.JSONFalse)
input.Restart = false
}
err := self.StartGuestDeployTask(ctx, userCred, kwargs, "deploy", "")
err := self.StartGuestDeployTask(ctx, userCred, jsonutils.Marshal(input).(*jsonutils.JSONDict), "deploy", "")
if err != nil {
return nil, err
}
@@ -785,33 +761,24 @@ 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) {
diskId, err := data.GetString("disk_id")
if err != nil {
log.Errorln(err)
func (self *SGuest) PerformAttachdisk(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerAttachDiskInput) (jsonutils.JSONObject, error) {
if len(input.DiskId) == 0 {
return nil, httperrors.NewMissingParameterError("disk_id")
}
disk, err := DiskManager.FetchByIdOrName(userCred, diskId)
if err != nil && err != sql.ErrNoRows {
log.Errorln(err)
return nil, err
}
if disk == nil {
return nil, httperrors.NewResourceNotFoundError("Disk %s not found", diskId)
}
if err := self.ValidateAttachDisk(ctx, disk.(*SDisk)); err != nil {
diskObj, err := validators.ValidateModel(userCred, DiskManager, &input.DiskId)
if err != nil {
return nil, err
}
taskData := data.(*jsonutils.JSONDict)
taskData.Set("disk_id", jsonutils.NewString(disk.GetId()))
if err := self.ValidateAttachDisk(ctx, diskObj.(*SDisk)); err != nil {
return nil, err
}
taskData := jsonutils.NewDict()
taskData.Add(jsonutils.NewString(input.DiskId), "disk_id")
self.SetStatus(userCred, api.VM_ATTACH_DISK, "")
if err := self.GetDriver().StartGuestAttachDiskTask(ctx, userCred, self, taskData, ""); err != nil {
return nil, err
}
return nil, nil
return nil, self.GetDriver().StartGuestAttachDiskTask(ctx, userCred, self, taskData, "")
}
func (self *SGuest) StartSyncTask(ctx context.Context, userCred mcclient.TokenCredential, firewallOnly bool,
@@ -1378,25 +1345,22 @@ func (self *SGuest) PerformAssignSecgroup(ctx context.Context, userCred mcclient
return nil, httperrors.NewMissingParameterError("secgroup_id")
}
secgroup, err := SecurityGroupManager.FetchByIdOrName(userCred, input.SecgroupId)
if err != nil {
if errors.Cause(err) == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError2("secgroup", input.SecgroupId)
}
return nil, httperrors.NewGeneralError(errors.Wrapf(err, "SecurityGroupManager.FetchByIdOrName(%s)", input.SecgroupId))
}
err = SecurityGroupManager.ValidateName(secgroup.GetName())
if err != nil {
return nil, httperrors.NewInputParameterError("The secgroup name %s does not meet the requirements, please change the name", secgroup.GetName())
}
err = self.saveDefaultSecgroupId(userCred, secgroup.GetId())
secObj, err := validators.ValidateModel(userCred, SecurityGroupManager, &input.SecgroupId)
if err != nil {
return nil, err
}
notes := map[string]string{"name": secgroup.GetName(), "id": secgroup.GetId()}
err = SecurityGroupManager.ValidateName(secObj.GetName())
if err != nil {
return nil, httperrors.NewInputParameterError("The secgroup name %s does not meet the requirements, please change the name", secObj.GetName())
}
err = self.saveDefaultSecgroupId(userCred, input.SecgroupId)
if err != nil {
return nil, err
}
notes := map[string]string{"name": secObj.GetName(), "id": secObj.GetId()}
logclient.AddActionLogWithContext(ctx, self, logclient.ACT_VM_ASSIGNSECGROUP, notes, userCred, true)
return nil, self.StartSyncTask(ctx, userCred, true, "")
}
@@ -1554,12 +1518,10 @@ func (self *SGuest) PerformRebuildRoot(ctx context.Context, userCred mcclient.To
return nil, err
}
imageId := input.GetImageName()
if len(imageId) > 0 {
img, err := CachedimageManager.getImageInfo(ctx, userCred, imageId, false)
if len(input.ImageId) > 0 {
img, err := CachedimageManager.getImageInfo(ctx, userCred, input.ImageId, false)
if err != nil {
return nil, httperrors.NewNotFoundError("failed to find %s", imageId)
return nil, httperrors.NewNotFoundError("failed to find %s", input.ImageId)
}
err = self.GetDriver().ValidateImage(ctx, img)
if err != nil {
@@ -1587,12 +1549,12 @@ func (self *SGuest) PerformRebuildRoot(ctx context.Context, userCred mcclient.To
if len(osName) == 0 && len(osType) == 0 && strings.ToLower(osType) != strings.ToLower(osName) {
return nil, httperrors.NewBadRequestError("Cannot switch OS between %s-%s", osName, osType)
}
imageId = img.Id
input.ImageId = img.Id
}
templateId := self.GetTemplateId()
if templateId != imageId && len(templateId) > 0 && len(imageId) > 0 && !self.GetDriver().IsRebuildRootSupportChangeUEFI() {
q := CachedimageManager.Query().In("id", []string{imageId, templateId})
if templateId != input.ImageId && len(templateId) > 0 && len(input.ImageId) > 0 && !self.GetDriver().IsRebuildRootSupportChangeUEFI() {
q := CachedimageManager.Query().In("id", []string{input.ImageId, templateId})
images := []SCachedimage{}
err := db.FetchModelObjects(CachedimageManager, q, &images)
if err != nil {
@@ -1608,11 +1570,11 @@ func (self *SGuest) PerformRebuildRoot(ctx context.Context, userCred mcclient.To
return nil, httperrors.NewInputParameterError("%v", err)
}
if !self.GetDriver().IsRebuildRootSupportChangeImage() && len(imageId) > 0 {
if !self.GetDriver().IsRebuildRootSupportChangeImage() && len(input.ImageId) > 0 {
if len(templateId) == 0 {
return nil, httperrors.NewBadRequestError("No template for root disk, cannot rebuild root")
}
if imageId != templateId {
if input.ImageId != templateId {
return nil, httperrors.NewInputParameterError("%s not support rebuild root with a different image", self.GetDriver().GetHypervisor())
}
}
@@ -1641,18 +1603,13 @@ func (self *SGuest) PerformRebuildRoot(ctx context.Context, userCred mcclient.To
}
}
keypairStr := input.GetKeypairName()
if len(keypairStr) > 0 {
keypairObj, err := KeypairManager.FetchByIdOrName(userCred, keypairStr)
if len(input.KeypairId) > 0 {
_, err := validators.ValidateModel(userCred, KeypairManager, &input.KeypairId)
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError("keypair %s not found", keypairStr)
} else {
return nil, httperrors.NewGeneralError(err)
}
return nil, err
}
if self.KeypairId != keypairObj.GetId() {
err = self.setKeypairId(userCred, keypairObj.GetId())
if self.KeypairId != input.KeypairId {
err = self.setKeypairId(userCred, input.KeypairId)
if err != nil {
return nil, httperrors.NewGeneralError(err)
}
@@ -1670,7 +1627,7 @@ func (self *SGuest) PerformRebuildRoot(ctx context.Context, userCred mcclient.To
allDisks = *input.AllDisks
}
return nil, self.StartRebuildRootTask(ctx, userCred, imageId, needStop, autoStart, passwd, resetPasswd, allDisks)
return nil, self.StartRebuildRootTask(ctx, userCred, input.ImageId, needStop, autoStart, passwd, resetPasswd, allDisks)
}
func (self *SGuest) GetTemplateId() string {
@@ -1819,54 +1776,44 @@ func (self *SGuest) AllowPerformDetachdisk(ctx context.Context, userCred mcclien
return self.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, self, "detachdisk")
}
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 {
func (self *SGuest) PerformDetachdisk(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerDetachDiskInput) (jsonutils.JSONObject, error) {
if len(input.DiskId) == 0 {
return nil, httperrors.NewMissingParameterError("disk_id")
}
keepDisk := jsonutils.QueryBoolean(data, "keep_disk", false)
iDisk, err := DiskManager.FetchByIdOrName(userCred, diskId)
diskObj, err := validators.ValidateModel(userCred, DiskManager, &input.DiskId)
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewNotFoundError("failed to find disk %s", diskId)
}
return nil, httperrors.NewGeneralError(err)
return nil, err
}
disk := iDisk.(*SDisk)
if disk != nil {
attached, err := self.isAttach2Disk(disk)
if err != nil {
return nil, httperrors.NewInternalServerError("check isAttach2Disk fail %s", err)
}
if attached {
if disk.DiskType == api.DISK_TYPE_SYS {
return nil, httperrors.NewUnsupportOperationError("Cannot detach sys 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")
}
err = self.GetDriver().ValidateDetachDisk(ctx, userCred, self, disk)
if err != nil {
return nil, err
}
if utils.IsInStringArray(self.Status, detachDiskStatus) {
self.SetStatus(userCred, api.VM_DETACH_DISK, "")
err = self.StartGuestDetachdiskTask(ctx, userCred, disk, keepDisk, "", false)
return nil, err
} 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)
}
disk := diskObj.(*SDisk)
attached, err := self.isAttach2Disk(disk)
if err != nil {
return nil, httperrors.NewInternalServerError("check isAttach2Disk fail %s", err)
}
return nil, httperrors.NewResourceNotFoundError("Disk %s not found", diskId)
if !attached {
return nil, nil
}
if disk.DiskType == api.DISK_TYPE_SYS {
return nil, httperrors.NewUnsupportOperationError("Cannot detach sys disk")
}
detachDiskStatus, err := self.GetDriver().GetDetachDiskStatus()
if err != nil {
return nil, err
}
if input.KeepDisk && !self.GetDriver().CanKeepDetachDisk() {
return nil, httperrors.NewInputParameterError("Cannot keep detached disk")
}
err = self.GetDriver().ValidateDetachDisk(ctx, userCred, self, disk)
if err != nil {
return nil, err
}
if utils.IsInStringArray(self.Status, detachDiskStatus) {
self.SetStatus(userCred, api.VM_DETACH_DISK, "")
err = self.StartGuestDetachdiskTask(ctx, userCred, disk, input.KeepDisk, "", false)
return nil, err
}
return nil, httperrors.NewInvalidStatusError("Server in %s not able to detach disk", self.Status)
}
func (self *SGuest) StartGuestDetachdiskTask(
@@ -2501,7 +2448,7 @@ func (self *SGuest) AllowPerformChangeConfig(ctx context.Context, userCred mccli
return self.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, self, "change-config")
}
func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerChangeConfigInput) (jsonutils.JSONObject, error) {
if !self.GetDriver().AllowReconfigGuest() {
return nil, httperrors.NewInvalidStatusError("Not allow to change config")
}
@@ -2518,18 +2465,17 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
return nil, httperrors.NewInvalidStatusError("Cannot change config in %s", self.Status)
}
host, _ := self.GetHost()
if host == nil {
return nil, httperrors.NewInvalidStatusError("No valid host")
_, err = self.GetHost()
if err != nil {
return nil, errors.Wrapf(err, "GetHost")
}
var addCpu, addMem int
var cpuChanged, memChanged bool
confs := jsonutils.NewDict()
skuId := jsonutils.GetAnyString(data, []string{"instance_type", "sku", "flavor"})
if len(skuId) > 0 {
sku, err := ServerSkuManager.FetchSkuByNameAndProvider(skuId, self.GetDriver().GetProvider(), true)
if len(input.InstanceType) > 0 {
sku, err := ServerSkuManager.FetchSkuByNameAndProvider(input.InstanceType, self.GetDriver().GetProvider(), true)
if err != nil {
return nil, err
}
@@ -2552,40 +2498,25 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
addMem = sku.MemorySizeMB - self.VmemSize
}
}
} else {
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")
}
if nVcpu != int64(self.VcpuCount) {
cpuChanged = true
addCpu = int(nVcpu - int64(self.VcpuCount))
err = confs.Add(jsonutils.NewInt(nVcpu), "vcpu_count")
if err != nil {
return nil, httperrors.NewBadRequestError("Params vcpu_count parse error")
}
}
if input.VcpuCount != self.VcpuCount {
cpuChanged = true
addCpu = input.VcpuCount - self.VcpuCount
confs.Add(jsonutils.NewInt(int64(input.VcpuCount)), "vcpu_count")
}
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 !regutils.MatchSize(input.VmemSize) {
return nil, httperrors.NewBadRequestError("Memory size must be number[+unit], like 256M, 1G or 256")
}
nVmem, err := fileutils.GetSizeMb(input.VmemSize, 'M', 1024)
if err != nil {
httperrors.NewBadRequestError("Params vmem_size parse error")
}
if nVmem != self.VmemSize {
memChanged = true
addMem = nVmem - self.VmemSize
err = confs.Add(jsonutils.NewInt(int64(nVmem)), "vmem_size")
if err != nil {
httperrors.NewBadRequestError("Params vmem_size parse error")
}
if nVmem != self.VmemSize {
memChanged = true
addMem = nVmem - self.VmemSize
err = confs.Add(jsonutils.NewInt(int64(nVmem)), "vmem_size")
if err != nil {
return nil, httperrors.NewBadRequestError("Params vmem_size parse error")
}
return nil, httperrors.NewBadRequestError("Params vmem_size parse error")
}
}
}
@@ -2607,42 +2538,33 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
var newDisks = make([]*api.DiskConfig, 0)
var resizeDisks = jsonutils.NewArray()
var inputDisks = make([]*api.DiskConfig, 0)
if disksConf, err := data.Get("disks"); err == nil {
if err = disksConf.Unmarshal(&inputDisks); err != nil {
return nil, httperrors.NewInputParameterError("Unmarshal disks configure error %s", err)
}
}
var schedInputDisks = make([]*api.DiskConfig, 0)
var diskIdx = 1
for _, diskConf := range inputDisks {
diskConf, err = parseDiskInfo(ctx, userCred, diskConf)
if err != nil {
return nil, httperrors.NewBadRequestError("Parse disk info error: %s", err)
for i := range input.Disks {
disk := input.Disks[i]
if len(disk.Backend) == 0 {
disk.Backend = self.getDefaultStorageType()
}
if len(diskConf.Backend) == 0 {
diskConf.Backend = self.getDefaultStorageType()
}
if diskConf.SizeMb > 0 {
if disk.SizeMb > 0 {
if diskIdx >= len(disks) {
newDisks = append(newDisks, diskConf)
newDisks = append(newDisks, &disk)
newDiskIdx += 1
addDisk += diskConf.SizeMb
schedInputDisks = append(schedInputDisks, diskConf)
addDisk += disk.SizeMb
schedInputDisks = append(schedInputDisks, &disk)
} else {
disk := disks[diskIdx].GetDisk()
oldSize := disk.DiskSize
if diskConf.SizeMb < oldSize {
gDisk := disks[diskIdx].GetDisk()
oldSize := gDisk.DiskSize
if disk.SizeMb < oldSize {
return nil, httperrors.NewInputParameterError("Cannot reduce disk size")
} else if diskConf.SizeMb > oldSize {
arr := jsonutils.NewArray(jsonutils.NewString(disks[diskIdx].DiskId), jsonutils.NewInt(int64(diskConf.SizeMb)))
}
if disk.SizeMb > oldSize {
arr := jsonutils.NewArray(jsonutils.NewString(disks[diskIdx].DiskId), jsonutils.NewInt(int64(disk.SizeMb)))
resizeDisks.Add(arr)
addDisk += diskConf.SizeMb - oldSize
addDisk += disk.SizeMb - oldSize
storage, _ := disks[diskIdx].GetDisk().GetStorage()
schedInputDisks = append(schedInputDisks, &api.DiskConfig{
SizeMb: addDisk,
Index: diskConf.Index,
Index: disk.Index,
Storage: storage.Id,
})
}
@@ -2654,7 +2576,7 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T
if resizeDisks.Length() > 0 {
confs.Add(resizeDisks, "resize")
}
if self.Status != api.VM_RUNNING && jsonutils.QueryBoolean(data, "auto_start", false) {
if self.Status != api.VM_RUNNING && input.AutoStart {
confs.Add(jsonutils.NewBool(true), "auto_start")
}
if self.Status == api.VM_RUNNING {
@@ -3250,23 +3172,19 @@ func (self *SGuest) setUserData(ctx context.Context, userCred mcclient.TokenCred
}
func (self *SGuest) AllowPerformUserData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return self.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, self, "userdata")
return self.IsOwner(userCred) || db.IsAdminAllowPerform(userCred, self, "user-data")
}
func (self *SGuest) PerformUserData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
userData, err := data.GetString("user_data")
if err != nil {
func (self *SGuest) PerformUserData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerUserDataInput) (jsonutils.JSONObject, error) {
if len(input.UserData) == 0 {
return nil, httperrors.NewMissingParameterError("user_data")
}
err = self.setUserData(ctx, userCred, userData)
err := self.setUserData(ctx, userCred, input.UserData)
if err != nil {
return nil, httperrors.NewGeneralError(err)
}
if len(self.HostId) > 0 {
err = self.StartSyncTask(ctx, userCred, false, "")
if err != nil {
return nil, httperrors.NewGeneralError(err)
}
return nil, self.StartSyncTask(ctx, userCred, false, "")
}
return nil, nil
}
+1 -3
View File
@@ -63,8 +63,6 @@ type IGuestDriver interface {
ValidateImage(ctx context.Context, image *cloudprovider.SImage) error
ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, data *api.ServerCreateInput) (*api.ServerCreateInput, error)
ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error)
ValidateCreateDataOnHost(ctx context.Context, userCred mcclient.TokenCredential, bmName string, host *SHost, input *api.ServerCreateInput) (*api.ServerCreateInput, error)
PrepareDiskRaidConfig(userCred mcclient.TokenCredential, host *SHost, params []*api.BaremetalDiskConfig, disks []*api.DiskConfig) ([]*api.DiskConfig, error)
@@ -109,7 +107,7 @@ type IGuestDriver interface {
RequestSyncstatusOnHost(ctx context.Context, guest *SGuest, host *SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error)
RequestStartOnHost(ctx context.Context, guest *SGuest, host *SHost, userCred mcclient.TokenCredential, task taskman.ITask) (jsonutils.JSONObject, error)
RequestStartOnHost(ctx context.Context, guest *SGuest, host *SHost, userCred mcclient.TokenCredential, task taskman.ITask) error
RequestStopOnHost(ctx context.Context, guest *SGuest, host *SHost, task taskman.ITask) error
+14 -88
View File
@@ -121,9 +121,6 @@ type SGuest struct {
// 秘钥对Id
KeypairId string `width:"36" charset:"ascii" nullable:"true" list:"user" create:"optional"`
// 宿主机Id
//HostId string `width:"36" charset:"ascii" nullable:"true" list:"admin" get:"admin" index:"true"`
// 备份机所在宿主机Id
BackupHostId string `width:"36" charset:"ascii" nullable:"true" list:"user" get:"user"`
@@ -132,7 +129,7 @@ type SGuest struct {
Machine string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"`
Bios string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"`
// 操作系统类型
OsType string `width:"36" charset:"ascii" nullable:"true" list:"user" update:"user" create:"optional"`
OsType string `width:"36" charset:"ascii" nullable:"true" list:"user" create:"optional"`
FlavorId string `width:"36" charset:"ascii" nullable:"true" list:"user" create:"optional"`
@@ -962,73 +959,17 @@ func ValidateMemCpuData(vmemSize, vcpuCount int, hypervisor string) (int, int, e
return vmemSize, vcpuCount, nil
}
func (self *SGuest) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
func (self *SGuest) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerUpdateInput) (api.ServerUpdateInput, error) {
if len(input.Name) > 0 && len(input.Name) < 2 {
return input, httperrors.NewInputParameterError("name is too short")
}
var err error
var vmemSize int
var vcpuCount int
driver := GetDriver(self.Hypervisor)
if memSize, _ := data.Int("vmem_size"); memSize != 0 {
vmemSize, err = ValidateMemData(int(memSize), driver)
if err != nil {
return nil, err
}
}
if cpuCount, _ := data.Int("vcpu_count"); cpuCount != 0 {
vcpuCount, err = ValidateCpuData(int(cpuCount), driver)
if err != nil {
return nil, err
}
}
if vmemSize > 0 || vcpuCount > 0 {
if !utils.IsInStringArray(self.Status, []string{api.VM_READY}) && self.GetHypervisor() != api.HYPERVISOR_CONTAINER {
return nil, httperrors.NewInvalidStatusError("Cannot modify Memory and CPU in status %s", self.Status)
}
if self.GetHypervisor() == api.HYPERVISOR_BAREMETAL {
return nil, httperrors.NewInputParameterError("Cannot modify memory for baremetal")
}
}
if vmemSize > 0 {
data.Add(jsonutils.NewInt(int64(vmemSize)), "vmem_size")
}
if vcpuCount > 0 {
data.Add(jsonutils.NewInt(int64(vcpuCount)), "vcpu_count")
}
data, err = self.GetDriver().ValidateUpdateData(ctx, userCred, data)
input.VirtualResourceBaseUpdateInput, err = self.SVirtualResourceBase.ValidateUpdateData(ctx, userCred, query, input.VirtualResourceBaseUpdateInput)
if err != nil {
return nil, err
return input, errors.Wrap(err, "SVirtualResourceBase.ValidateUpdateData")
}
if vcpuCount > 0 || vmemSize > 0 {
quota, err := self.checkUpdateQuota(ctx, userCred, vcpuCount, vmemSize)
if err != nil {
return nil, httperrors.NewOutOfQuotaError("%v", err)
}
if !quota.IsEmpty() {
data.Add(jsonutils.Marshal(quota), "pending_usage")
}
}
if data.Contains("name") {
if name, _ := data.GetString("name"); len(name) < 2 {
return nil, httperrors.NewInputParameterError("name is too short")
}
}
input := apis.VirtualResourceBaseUpdateInput{}
err = data.Unmarshal(&input)
if err != nil {
return nil, errors.Wrap(err, "data.Unmarshal")
}
input, err = self.SVirtualResourceBase.ValidateUpdateData(ctx, userCred, query, input)
if err != nil {
return nil, errors.Wrap(err, "SVirtualResourceBase.ValidateUpdateData")
}
data.Update(jsonutils.Marshal(input))
return data, nil
return input, nil
}
func serverCreateInput2ComputeQuotaKeys(input api.ServerCreateInput, ownerId mcclient.IIdentityProvider) SComputeResourceKeys {
@@ -1667,15 +1608,6 @@ func (manager *SGuestManager) validateEip(userCred mcclient.TokenCredential, inp
func (self *SGuest) PostUpdate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) {
self.SVirtualResourceBase.PostUpdate(ctx, userCred, query, data)
if data.Contains("pending_usage") {
quota := SQuota{}
data.Unmarshal(&quota, "pending_usage")
quotas.CancelPendingUsage(ctx, userCred, &quota, &quota, true)
}
self.StartSyncTask(ctx, userCred, true, "")
if data.Contains("name") || data.Contains("__meta__") {
err := self.StartRemoteUpdateTask(ctx, userCred, false, "")
if err != nil {
@@ -5023,21 +4955,15 @@ func (self *SGuest) SyncVMSecgroups(ctx context.Context, userCred mcclient.Token
func (self *SGuest) GetIVM() (cloudprovider.ICloudVM, error) {
if len(self.ExternalId) == 0 {
msg := fmt.Sprintf("GetIVM: not managed by a provider")
log.Errorf(msg)
return nil, fmt.Errorf(msg)
return nil, errors.Wrapf(cloudprovider.ErrNotFound, "empty externalId")
}
host, _ := self.GetHost()
if host == nil {
msg := fmt.Sprintf("GetIVM: No valid host")
log.Errorf(msg)
return nil, fmt.Errorf(msg)
host, err := self.GetHost()
if err != nil {
return nil, errors.Wrapf(err, "GetHost")
}
ihost, err := host.GetIHost()
if err != nil {
msg := fmt.Sprintf("GetIVM: getihost fail %s", err)
log.Errorf(msg)
return nil, fmt.Errorf(msg)
return nil, errors.Wrapf(err, "GetIHost")
}
return ihost.GetIVMById(self.ExternalId)
}
+43 -30
View File
@@ -54,6 +54,7 @@ import (
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudcommon/policy"
"yunion.io/x/onecloud/pkg/cloudcommon/validators"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/options"
"yunion.io/x/onecloud/pkg/httperrors"
@@ -1434,29 +1435,25 @@ func (manager *SNetworkManager) validateEnsureWire(ctx context.Context, userCred
return
}
func (manager *SNetworkManager) validateEnsureZoneVpc(ctx context.Context, userCred mcclient.TokenCredential, input api.NetworkCreateInput) (w *SWire, v *SVpc, cr *SCloudregion, err error) {
defer func() {
if cause := errors.Cause(err); cause == sql.ErrNoRows {
err = httperrors.NewResourceNotFoundError("%s", err)
}
}()
zObj, err := ZoneManager.FetchByIdOrName(userCred, input.Zone)
func (manager *SNetworkManager) validateEnsureZoneVpc(ctx context.Context, userCred mcclient.TokenCredential, input api.NetworkCreateInput) (*SWire, *SVpc, *SCloudregion, error) {
zObj, err := validators.ValidateModel(userCred, ZoneManager, &input.Zone)
if err != nil {
err = errors.Wrapf(err, "zone %s", input.Zone)
return
return nil, nil, nil, err
}
z := zObj.(*SZone)
vObj, err := VpcManager.FetchByIdOrName(userCred, input.Vpc)
vObj, err := validators.ValidateModel(userCred, VpcManager, &input.Vpc)
if err != nil {
err = errors.Wrapf(err, "vpc %s", input.Vpc)
return
return nil, nil, nil, err
}
v = vObj.(*SVpc)
v := vObj.(*SVpc)
var wires []SWire
cr, err := z.GetRegion()
if err != nil {
return nil, nil, nil, err
}
// 华为云,ucloud wire zone_id 为空
cr, _ = z.GetRegion()
var wires []SWire
if utils.IsInStringArray(cr.Provider, api.REGIONAL_NETWORK_PROVIDERS) {
wires, err = WireManager.getWiresByVpcAndZone(v, nil)
} else {
@@ -1464,25 +1461,41 @@ func (manager *SNetworkManager) validateEnsureZoneVpc(ctx context.Context, userC
}
if err != nil {
return
} else if len(wires) > 1 {
err = httperrors.NewConflictError("found %d wires for zone %s and vpc %s", len(wires), input.Zone, input.Vpc)
return
} else if len(wires) == 1 {
w = &wires[0]
return
return nil, nil, nil, err
}
// wire not found. We auto create one for OneCloud vpc
if cr.Provider == api.CLOUD_PROVIDER_ONECLOUD {
w, err = v.initWire(ctx, z)
if len(wires) > 1 {
return nil, nil, nil, httperrors.NewConflictError("found %d wires for zone %s and vpc %s", len(wires), input.Zone, input.Vpc)
}
if len(wires) == 1 {
return &wires[0], v, cr, nil
}
externalId := ""
if cr.Provider == api.CLOUD_PROVIDER_CLOUDPODS {
iVpc, err := v.GetIVpc()
if err != nil {
err = errors.Wrapf(err, "vpc %s init wire", v.Id)
return
return nil, nil, nil, err
}
return
iWire, err := iVpc.CreateIWire(&cloudprovider.SWireCreateOptions{
Name: fmt.Sprintf("vpc-%s", v.Name),
ZoneId: z.ExternalId,
Bandwidth: 10000,
Mtu: 1500,
})
if err != nil {
return nil, nil, nil, errors.Wrapf(err, "CreateIWire")
}
externalId = iWire.GetGlobalId()
}
err = httperrors.NewNotFoundError("wire not found for zone %s and vpc %s", input.Zone, input.Vpc)
return
// wire not found. We auto create one for OneCloud vpc
if cr.Provider == api.CLOUD_PROVIDER_ONECLOUD || cr.Provider == api.CLOUD_PROVIDER_CLOUDPODS {
w, err := v.initWire(ctx, z, externalId)
if err != nil {
return nil, nil, nil, errors.Wrapf(err, "vpc %s init wire", v.Id)
}
return w, v, cr, nil
}
return nil, nil, nil, httperrors.NewNotFoundError("wire not found for zone %s and vpc %s", input.Zone, input.Vpc)
}
func (manager *SNetworkManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, input api.NetworkCreateInput) (api.NetworkCreateInput, error) {
+2 -1
View File
@@ -1285,7 +1285,7 @@ func (vpc *SVpc) PerformSync(ctx context.Context, userCred mcclient.TokenCredent
return nil, httperrors.NewUnsupportOperationError("on-premise vpc cannot sync status")
}
func (self *SVpc) initWire(ctx context.Context, zone *SZone) (*SWire, error) {
func (self *SVpc) initWire(ctx context.Context, zone *SZone, externalId string) (*SWire, error) {
wire := &SWire{
Bandwidth: 10000,
Mtu: 1500,
@@ -1293,6 +1293,7 @@ func (self *SVpc) initWire(ctx context.Context, zone *SZone) (*SWire, error) {
wire.VpcId = self.Id
wire.ZoneId = zone.Id
wire.IsEmulated = true
wire.ExternalId = externalId
wire.Name = fmt.Sprintf("vpc-%s", self.Name)
wire.DomainId = self.DomainId
+2 -1
View File
@@ -222,7 +222,8 @@ func (self *GuestStartAndSyncToBackupTask) checkTemplete(ctx context.Context, gu
func (self *GuestStartAndSyncToBackupTask) OnCheckTemplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) {
self.SetStage("OnStartBackupGuest", nil)
host := models.HostManager.FetchHostById(guest.BackupHostId)
if _, err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self); err != nil {
err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self)
if err != nil {
self.SetStageFailed(ctx, jsonutils.NewString(err.Error()))
}
}
+4 -19
View File
@@ -45,34 +45,19 @@ func (self *GuestStartTask) RequestStart(ctx context.Context, guest *models.SGue
self.SetStage("OnStartComplete", nil)
host, _ := guest.GetHost()
guest.SetStatus(self.UserCred, api.VM_STARTING, "")
result, err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self)
err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self)
if err != nil {
self.OnStartCompleteFailed(ctx, guest, jsonutils.NewString(err.Error()))
} else {
if result != nil && jsonutils.QueryBoolean(result, "is_running", false) {
// guest.SetStatus(self.UserCred, models.VM_RUNNING, "start")
// self.taskComplete(ctx, guest)
self.OnStartComplete(ctx, guest, nil)
}
return
}
self.OnStartComplete(ctx, guest, nil)
}
func (self *GuestStartTask) OnStartComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
db.OpsLog.LogEvent(guest, db.ACT_START, guest.GetShortDesc(ctx), self.UserCred)
self.SetStage("OnGuestSyncstatusAfterStart", nil)
if guest.Hypervisor != api.HYPERVISOR_KVM {
guest.StartSyncstatus(ctx, self.UserCred, self.GetTaskId())
} else {
logclient.AddActionLogWithStartable(self, guest, logclient.ACT_VM_START, "success", self.UserCred, true)
self.taskComplete(ctx, guest)
}
}
func (self *GuestStartTask) OnGuestSyncstatusAfterStart(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
guest := obj.(*models.SGuest)
logclient.AddActionLogWithStartable(self, guest, logclient.ACT_VM_START, "success", self.UserCred, true)
self.taskComplete(ctx, guest)
logclient.AddActionLogWithStartable(self, guest, logclient.ACT_VM_START, "", self.UserCred, true)
}
func (self *GuestStartTask) OnStartCompleteFailed(ctx context.Context, obj db.IStandaloneModel, err jsonutils.JSONObject) {
+8 -12
View File
@@ -18,7 +18,7 @@ import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/onecloud/pkg/apis"
api "yunion.io/x/onecloud/pkg/apis/compute"
@@ -44,34 +44,30 @@ func (self *GuestStopTask) OnInit(ctx context.Context, obj db.IStandaloneModel,
}
func (self *GuestStopTask) stopGuest(ctx context.Context, guest *models.SGuest) {
host, _ := guest.GetHost()
if host == nil {
self.OnGuestStopTaskCompleteFailed(ctx, guest, jsonutils.NewString("no associated host"))
host, err := guest.GetHost()
if err != nil {
self.OnGuestStopTaskCompleteFailed(ctx, guest, jsonutils.NewString(errors.Wrapf(err, "GetHost").Error()))
return
}
if !self.IsSubtask() {
guest.SetStatus(self.UserCred, api.VM_STOPPING, "")
guest.SetStatus(self.GetUserCred(), api.VM_STOPPING, "")
}
self.SetStage("OnGuestStopTaskComplete", nil)
err := guest.GetDriver().RequestStopOnHost(ctx, guest, host, self)
err = guest.GetDriver().RequestStopOnHost(ctx, guest, host, self)
if err != nil {
log.Errorf("RequestStopOnHost fail %s", err)
self.OnGuestStopTaskCompleteFailed(ctx, guest, jsonutils.NewString(err.Error()))
}
}
func (self *GuestStopTask) OnGuestStopTaskComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) {
if !self.IsSubtask() {
guest.StartSyncstatus(ctx, self.UserCred, "")
// guest.SetStatus(self.UserCred, api.VM_READY, "")
}
db.OpsLog.LogEvent(guest, db.ACT_STOP, guest.GetShortDesc(ctx), self.UserCred)
models.HostManager.ClearSchedDescCache(guest.HostId)
logclient.AddActionLogWithStartable(self, guest, logclient.ACT_VM_STOP, "success", self.UserCred, true)
self.SetStageComplete(ctx, nil)
if guest.Status == api.VM_READY && guest.DisableDelete.IsFalse() && guest.ShutdownBehavior == api.SHUTDOWN_TERMINATE {
guest.StartAutoDeleteGuestTask(ctx, self.UserCred, "")
return
}
logclient.AddActionLogWithStartable(self, guest, logclient.ACT_VM_STOP, "success", self.UserCred, true)
}
func (self *GuestStopTask) OnGuestStopTaskCompleteFailed(ctx context.Context, guest *models.SGuest, reason jsonutils.JSONObject) {
+2 -5
View File
@@ -47,13 +47,10 @@ func (self *HAGuestStartTask) RequestStartBacking(ctx context.Context, guest *mo
host := models.HostManager.FetchHostById(guest.BackupHostId)
guest.SetStatus(self.UserCred, api.VM_BACKUP_STARTING, "")
result, err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self)
err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self)
if err != nil {
self.OnStartCompleteFailed(ctx, guest, jsonutils.NewString(err.Error()))
} else {
if result != nil && jsonutils.QueryBoolean(result, "is_running", false) {
self.RequestStart(ctx, guest)
}
return
}
}
+4 -4
View File
@@ -633,17 +633,17 @@ func (o *ServerCancelDeleteOptions) Description() string {
type ServerDeployOptions struct {
ServerIdOptions
Keypair string `help:"ssh Keypair used for login" json:"-"`
DeleteKeypair *bool `help:"Remove ssh Keypairs" json:"-"`
DeleteKeypair bool `help:"Remove ssh Keypairs" json:"-"`
Deploy []string `help:"Specify deploy files in virtual server file system" json:"-"`
ResetPassword *bool `help:"Force reset password"`
ResetPassword bool `help:"Force reset password"`
Password string `help:"Default user password"`
AutoStart *bool `help:"Auto start server after deployed"`
AutoStart bool `help:"Auto start server after deployed"`
}
func (opts *ServerDeployOptions) Params() (jsonutils.JSONObject, error) {
params := new(computeapi.ServerDeployInput)
{
if opts.DeleteKeypair != nil {
if opts.DeleteKeypair == true {
params.DeleteKeypair = opts.DeleteKeypair
} else if len(opts.Keypair) > 0 {
params.Keypair = opts.Keypair
+13 -6
View File
@@ -49,6 +49,7 @@ type ModelManager interface {
Delete(session *mcclient.ClientSession, id string, param jsonutils.JSONObject) (jsonutils.JSONObject, error)
PerformAction(session *mcclient.ClientSession, id string, action string, params jsonutils.JSONObject) (jsonutils.JSONObject, error)
Get(session *mcclient.ClientSession, id string, params jsonutils.JSONObject) (jsonutils.JSONObject, error)
Update(session *mcclient.ClientSession, id string, params jsonutils.JSONObject) (jsonutils.JSONObject, error)
}
type CloudpodsClientConfig struct {
@@ -130,27 +131,31 @@ func (self *SCloudpodsClient) get(manager ModelManager, id string, params map[st
return resp.Unmarshal(retVal)
}
func (self *SCloudpodsClient) perform(manager ModelManager, id, action string, params map[string]interface{}) error {
_, err := manager.PerformAction(self.s, id, action, jsonutils.Marshal(params))
return err
func (self *SCloudpodsClient) perform(manager ModelManager, id, action string, params interface{}) (jsonutils.JSONObject, error) {
return manager.PerformAction(self.s, id, action, jsonutils.Marshal(params))
}
func (self *SCloudpodsClient) delete(manager ModelManager, id string) error {
if len(id) == 0 {
return nil
}
params := map[string]interface{}{}
params := map[string]interface{}{"override_pending_delete": true}
_, err := manager.Delete(self.s, id, jsonutils.Marshal(params))
return err
}
func (self *SCloudpodsClient) update(manager ModelManager, id string, params interface{}) error {
_, err := manager.Update(self.s, id, jsonutils.Marshal(params))
return err
}
func (self *SCloudpodsClient) GetAccountId() string {
return self.authURL
}
func (self *SCloudpodsClient) GetSubAccounts() ([]cloudprovider.SSubAccount, error) {
return []cloudprovider.SSubAccount{
cloudprovider.SSubAccount{
{
Name: self.cpcfg.Name,
Account: self.cpcfg.Account,
HealthStatus: api.CLOUD_PROVIDER_HEALTH_NORMAL,
@@ -179,7 +184,9 @@ func (self *SCloudpodsClient) list(manager ModelManager, params map[string]inter
params = map[string]interface{}{}
}
for k, v := range defaultParams {
params[k] = v
if _, ok := params[k]; !ok {
params[k] = v
}
}
ret := []jsonutils.JSONObject{}
for {
+12 -7
View File
@@ -16,6 +16,7 @@ package cloudpods
import (
"context"
"fmt"
"time"
"yunion.io/x/jsonutils"
@@ -155,15 +156,17 @@ func (self *SDisk) GetExtSnapshotPolicyIds() ([]string, error) {
}
func (self *SDisk) Resize(ctx context.Context, sizeMb int64) error {
var err error
params := map[string]interface{}{
"size_mb": sizeMb,
"disk_id": self.Id,
input := api.ServerResizeDiskInput{
DiskResizeInput: api.DiskResizeInput{},
}
input.DiskId = self.Id
input.Size = fmt.Sprintf("%dM", sizeMb)
if len(self.Guests) > 0 {
err = self.region.perform(&modules.Servers, self.Guests[0].Id, "resize-disk", params)
_, err := self.region.perform(&modules.Servers, self.Guests[0].Id, "resize-disk", input)
return err
}
err = self.region.perform(&modules.Disks, self.Id, "resize", params)
_, err := self.region.perform(&modules.Disks, self.Id, "resize", input.DiskResizeInput)
return err
}
@@ -231,7 +234,9 @@ func (self *SStorage) GetIDiskById(id string) (cloudprovider.ICloudDisk, error)
}
func (self *SStorage) CreateIDisk(opts *cloudprovider.DiskCreateConfig) (cloudprovider.ICloudDisk, error) {
input := api.DiskCreateInput{}
input := api.DiskCreateInput{
DiskConfig: &api.DiskConfig{},
}
input.Name = opts.Name
input.Description = opts.Desc
input.SizeMb = opts.SizeGb * 1024
+7 -6
View File
@@ -99,20 +99,21 @@ func (self *SEip) IsAutoRenew() bool {
}
func (self *SEip) Associate(opts *cloudprovider.AssociateConfig) error {
params := map[string]interface{}{
"associate_type": opts.AssociateType,
"instance_id": opts.InstanceId,
}
input := api.ElasticipAssociateInput{}
input.InstanceType = opts.AssociateType
input.InstanceId = opts.InstanceId
switch opts.AssociateType {
case api.EIP_ASSOCIATE_TYPE_SERVER:
default:
return cloudprovider.ErrNotImplemented
}
return self.region.perform(&modules.Elasticips, self.Id, "associate", params)
_, err := self.region.perform(&modules.Elasticips, self.Id, "associate", input)
return err
}
func (self *SEip) Dissociate() error {
return self.region.perform(&modules.Elasticips, self.Id, "dissociate", nil)
_, err := self.region.perform(&modules.Elasticips, self.Id, "dissociate", nil)
return err
}
func (self *SEip) GetCreatedAt() time.Time {
+152 -47
View File
@@ -16,13 +16,18 @@ package cloudpods
import (
"context"
"fmt"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
"yunion.io/x/onecloud/pkg/apis"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/mcclient/modules"
"yunion.io/x/onecloud/pkg/mcclient/modules/webconsole"
"yunion.io/x/onecloud/pkg/multicloud"
)
@@ -155,17 +160,17 @@ func (self *SInstance) GetProjectId() string {
}
func (self *SInstance) AssignSecurityGroup(id string) error {
params := map[string]interface{}{
"secgroup_id": id,
}
return self.host.zone.region.perform(&modules.Servers, self.Id, "assign-secgroup", params)
input := api.GuestAssignSecgroupInput{}
input.SecgroupId = id
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "assign-secgroup", input)
return err
}
func (self *SInstance) SetSecurityGroups(ids []string) error {
params := map[string]interface{}{
"secgroup_ids": ids,
}
return self.host.zone.region.perform(&modules.Servers, self.Id, "set-secgroup", params)
input := api.GuestSetSecgroupInput{}
input.SecgroupIds = ids
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "set-secgroup", input)
return err
}
func (self *SInstance) GetHypervisor() string {
@@ -173,64 +178,138 @@ func (self *SInstance) GetHypervisor() string {
}
func (self *SInstance) StartVM(ctx context.Context) error {
return self.host.zone.region.perform(&modules.Servers, self.Id, "start", nil)
if self.Status == api.VM_RUNNING {
return nil
}
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "start", nil)
return err
}
func (self *SInstance) StopVM(ctx context.Context, opts *cloudprovider.ServerStopOptions) error {
params := map[string]interface{}{
"is_force": opts.IsForce,
if self.Status == api.VM_READY {
return nil
}
return self.host.zone.region.perform(&modules.Servers, self.Id, "stop", params)
input := api.ServerStopInput{}
input.IsForce = opts.IsForce
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "stop", input)
return err
}
func (self *SInstance) DeleteVM(ctx context.Context) error {
if self.DisableDelete != nil && *self.DisableDelete {
input := api.ServerUpdateInput{}
disableDelete := false
input.DisableDelete = &disableDelete
self.host.zone.region.cli.update(&modules.Servers, self.Id, input)
}
return self.host.zone.region.cli.delete(&modules.Servers, self.Id)
}
func (self *SInstance) UpdateVM(ctx context.Context, name string) error {
return cloudprovider.ErrNotImplemented
if self.Name != name {
input := api.ServerUpdateInput{}
input.Name = name
self.host.zone.region.cli.update(&modules.Servers, self.Id, input)
return cloudprovider.WaitMultiStatus(self, []string{api.VM_READY, api.VM_RUNNING}, time.Second*5, time.Minute*3)
}
return nil
}
func (self *SInstance) UpdateUserData(userData string) error {
return cloudprovider.ErrNotImplemented
input := api.ServerUserDataInput{}
input.UserData = userData
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "user-data", input)
return err
}
func (self *SInstance) RebuildRoot(ctx context.Context, opts *cloudprovider.SManagedVMRebuildRootConfig) (string, error) {
params := map[string]interface{}{
"image_id": opts.ImageId,
"password": opts.Password,
input := api.ServerRebuildRootInput{}
input.ImageId = opts.ImageId
input.Password = opts.Password
if len(opts.PublicKey) > 0 {
keypairId, err := self.host.zone.region.syncKeypair(self.Name, opts.PublicKey)
if err != nil {
return "", errors.Wrapf(err, "syncKeypair")
}
input.KeypairId = keypairId
}
diskId := self.DisksInfo[0].Id
return diskId, self.host.zone.region.perform(&modules.Servers, self.Id, "rebuild-root", params)
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "rebuild-root", input)
if err != nil {
return "", err
}
return self.DisksInfo[0].Id, nil
}
func (self *SInstance) DeployVM(ctx context.Context, name string, username string, password string, publicKey string, deleteKeypair bool, description string) error {
params := map[string]interface{}{
"password": password,
input := api.ServerDeployInput{}
input.Password = password
if len(publicKey) > 0 {
keypairId, err := self.host.zone.region.syncKeypair(name, publicKey)
if err != nil {
return errors.Wrapf(err, "syncKeypair")
}
input.KeypairId = keypairId
}
return self.host.zone.region.perform(&modules.Servers, self.Id, "deploy", params)
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "deploy", input)
if err != nil {
return errors.Wrapf(err, "deploy")
}
return cloudprovider.WaitMultiStatus(self, []string{api.VM_READY, api.VM_RUNNING}, time.Second*5, time.Minute*3)
}
func (self *SInstance) ChangeConfig(ctx context.Context, config *cloudprovider.SManagedVMChangeConfig) error {
return cloudprovider.ErrNotImplemented
func (self *SInstance) ChangeConfig(ctx context.Context, opts *cloudprovider.SManagedVMChangeConfig) error {
input := api.ServerChangeConfigInput{}
input.VmemSize = fmt.Sprintf("%dM", opts.MemoryMB)
input.VcpuCount = opts.Cpu
input.InstanceType = opts.InstanceType
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "change-config", input)
return err
}
func (self *SInstance) GetVNCInfo() (jsonutils.JSONObject, error) {
return nil, cloudprovider.ErrNotImplemented
s := self.host.zone.region.cli.s
resp, err := webconsole.WebConsole.DoServerConnect(s, self.Id, nil)
if err != nil {
return nil, errors.Wrapf(err, "DoServerConnect")
}
data := struct {
ConnectParams string
ApiServer string
Session string
Protocol string
InstanceId string
InstanceName string
}{
Protocol: "cloudpods",
InstanceId: self.Id,
InstanceName: self.Name,
}
err = resp.Unmarshal(&data)
if err != nil {
return nil, errors.Wrapf(err, "resp.Unmarshal")
}
resp, err = modules.ServicesV3.GetSpecific(s, "common", "config", nil)
if err != nil {
return nil, errors.Wrapf(err, "GetSpecific")
}
data.ApiServer, _ = resp.GetString("config", "default", "api_server")
return jsonutils.Marshal(data), nil
}
func (self *SInstance) AttachDisk(ctx context.Context, diskId string) error {
params := map[string]interface{}{
"disk_id": diskId,
}
return self.host.zone.region.perform(&modules.Disks, self.Id, "attach-disk", params)
input := api.ServerAttachDiskInput{}
input.DiskId = diskId
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "attachdisk", input)
return err
}
func (self *SInstance) DetachDisk(ctx context.Context, diskId string) error {
params := map[string]interface{}{
"disk_id": diskId,
}
return self.host.zone.region.perform(&modules.Disks, self.Id, "detach-disk", params)
input := api.ServerDetachDiskInput{}
input.DiskId = diskId
input.KeepDisk = true
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "detachdisk", input)
return err
}
func (self *SInstance) CreateDisk(ctx context.Context, sizeMb int, uuid string, driver string) error {
@@ -238,21 +317,35 @@ func (self *SInstance) CreateDisk(ctx context.Context, sizeMb int, uuid string,
}
func (self *SInstance) MigrateVM(hostId string) error {
params := map[string]interface{}{
"prefer_host": hostId,
}
return self.host.zone.region.perform(&modules.Servers, self.Id, "migrate", params)
input := api.GuestMigrateInput{}
input.PreferHost = hostId
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "migrate", input)
return err
}
func (self *SInstance) LiveMigrateVM(hostId string) error {
params := map[string]interface{}{
"prefer_host": hostId,
}
return self.host.zone.region.perform(&modules.Servers, self.Id, "live-migrate", params)
input := api.GuestLiveMigrateInput{}
input.PreferHost = hostId
skipCpuCheck := true
input.SkipCpuCheck = &skipCpuCheck
_, err := self.host.zone.region.perform(&modules.Servers, self.Id, "live-migrate", input)
return err
}
func (self *SInstance) GetError() error {
return cloudprovider.ErrNotImplemented
if utils.IsInStringArray(self.Status, []string{api.VM_DISK_FAILED, api.VM_SCHEDULE_FAILED, api.VM_NETWORK_FAILED}) {
return fmt.Errorf("vm create failed with status %s", self.Status)
}
if self.Status == api.VM_DEPLOY_FAILED {
params := map[string]interface{}{"obj_id": self.Id, "success": false}
actions := []apis.OpsLogDetails{}
self.host.zone.region.list(&modules.Actions, params, &actions)
if len(actions) > 0 {
return fmt.Errorf(actions[0].Notes)
}
return fmt.Errorf("vm create failed with status %s", self.Status)
}
return nil
}
func (self *SInstance) CreateInstanceSnapshot(ctx context.Context, name string, desc string) (cloudprovider.ICloudInstanceSnapshot, error) {
@@ -272,7 +365,18 @@ func (self *SInstance) ResetToInstanceSnapshot(ctx context.Context, idStr string
}
func (self *SInstance) SaveImage(opts *cloudprovider.SaveImageOptions) (cloudprovider.ICloudImage, error) {
return nil, cloudprovider.ErrNotImplemented
input := api.ServerSaveImageInput{}
input.GenerateName = opts.Name
input.Notes = opts.Notes
resp, err := self.host.zone.region.perform(&modules.Servers, self.Id, "save-image", input)
if err != nil {
return nil, err
}
err = resp.Unmarshal(&input)
if err != nil {
return nil, err
}
return self.host.zone.region.GetImage(input.ImageId)
}
func (self *SInstance) AllocatePublicIpAddress() (string, error) {
@@ -316,14 +420,15 @@ func (self *SRegion) GetInstances(hostId string) ([]SInstance, error) {
}
func (self *SRegion) CreateInstance(hostId, hypervisor string, opts *cloudprovider.SManagedVMCreateConfig) (*SInstance, error) {
input := api.ServerCreateInput{}
input := api.ServerCreateInput{
ServerConfigs: &api.ServerConfigs{},
}
input.Name = opts.Name
input.Description = opts.Description
input.InstanceType = opts.InstanceType
input.VcpuCount = opts.Cpu
input.VmemSize = opts.MemoryMB
input.Password = opts.Password
input.LoginAccount = opts.Account
input.PublicIpBw = opts.PublicIpBw
input.PublicIpChargeType = string(opts.PublicIpChargeType)
input.ProjectId = opts.ProjectId
@@ -341,7 +446,7 @@ func (self *SRegion) CreateInstance(hostId, hypervisor string, opts *cloudprovid
input.Disks = append(input.Disks, &api.DiskConfig{
Index: 0,
ImageId: opts.ExternalImageId,
DiskType: "sys",
DiskType: api.DISK_TYPE_SYS,
SizeMb: opts.SysDisk.SizeGB * 1024,
Backend: opts.SysDisk.StorageType,
Storage: opts.SysDisk.StorageExternalId,
@@ -349,7 +454,7 @@ func (self *SRegion) CreateInstance(hostId, hypervisor string, opts *cloudprovid
for idx, disk := range opts.DataDisks {
input.Disks = append(input.Disks, &api.DiskConfig{
Index: idx + 1,
DiskType: "data",
DiskType: api.DISK_TYPE_DATA,
SizeMb: disk.SizeGB * 1024,
Backend: disk.StorageType,
Storage: disk.StorageExternalId,
+58
View File
@@ -0,0 +1,58 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package cloudpods
import (
"fmt"
"yunion.io/x/pkg/errors"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/mcclient/modules"
)
type SKeypair struct {
api.KeypairDetails
}
func (self *SRegion) GetKeypairs() ([]SKeypair, error) {
keypairs := []SKeypair{}
return keypairs, self.list(&modules.Keypairs, nil, &keypairs)
}
func (self *SRegion) CreateKeypair(name, publicKey string) (*SKeypair, error) {
input := api.KeypairCreateInput{}
input.GenerateName = fmt.Sprintf("keypair-for-%s", name)
input.PublicKey = publicKey
keypair := &SKeypair{}
return keypair, self.create(&modules.Keypairs, input, keypair)
}
func (self *SRegion) syncKeypair(serverName, publicKey string) (string, error) {
keypairs, err := self.GetKeypairs()
if err != nil {
return "", err
}
for _, keypair := range keypairs {
if keypair.PublicKey == publicKey {
return keypair.Id, nil
}
}
keypair, err := self.CreateKeypair(serverName, publicKey)
if err != nil {
return "", errors.Wrapf(err, "CreateKeypair")
}
return keypair.Id, nil
}
+3 -1
View File
@@ -17,6 +17,8 @@ package cloudpods
import (
"fmt"
"yunion.io/x/jsonutils"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/multicloud"
@@ -115,7 +117,7 @@ func (self *SRegion) create(manager ModelManager, params interface{}, retVal int
return self.cli.create(manager, params, retVal)
}
func (self *SRegion) perform(manager ModelManager, id, action string, params map[string]interface{}) error {
func (self *SRegion) perform(manager ModelManager, id, action string, params interface{}) (jsonutils.JSONObject, error) {
return self.cli.perform(manager, id, action, params)
}
+6 -2
View File
@@ -109,11 +109,15 @@ func (self *SZone) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
}
func (self *SZone) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
storage, err := self.region.GetStorage(id)
return self.region.GetIStorageById(id)
}
func (self *SRegion) GetIStorageById(id string) (cloudprovider.ICloudStorage, error) {
storage, err := self.GetStorage(id)
if err != nil {
return nil, err
}
storage.region = self.region
storage.region = self
return storage, nil
}
+29
View File
@@ -76,6 +76,10 @@ func (self *SVpc) GetRegion() cloudprovider.ICloudRegion {
return self.region
}
func (self *SVpc) GetExternalAccessMode() string {
return self.ExternalAccessMode
}
func (self *SVpc) Delete() error {
return self.region.cli.delete(&modules.Vpcs, self.Id)
}
@@ -93,6 +97,31 @@ func (self *SRegion) GetIVpcs() ([]cloudprovider.ICloudVpc, error) {
return ret, nil
}
func (self *SVpc) CreateIWire(opts *cloudprovider.SWireCreateOptions) (cloudprovider.ICloudWire, error) {
wire, err := self.region.CreateWire(opts, self.Id, self.DomainId, self.PublicScope, self.IsPublic)
if err != nil {
return nil, err
}
wire.vpc = self
return wire, nil
}
func (self *SRegion) CreateWire(opts *cloudprovider.SWireCreateOptions, vpcId, domainId, publicScope string, isPublic bool) (*SWire, error) {
input := api.WireCreateInput{}
input.GenerateName = opts.Name
input.Mtu = opts.Mtu
input.Bandwidth = opts.Bandwidth
input.VpcId = vpcId
input.DomainId = domainId
input.PublicScope = publicScope
input.IsPublic = &isPublic
input.ZoneId = opts.ZoneId
t := true
input.IsEmulated = &t
wire := &SWire{}
return wire, self.create(&modules.Wires, input, wire)
}
func (self *SRegion) CreateIVpc(name string, desc string, cidr string) (cloudprovider.ICloudVpc, error) {
input := api.VpcCreateInput{}
input.Name = name
+1 -1
View File
@@ -104,7 +104,7 @@ func (self *SVpc) GetIWires() ([]cloudprovider.ICloudWire, error) {
func (self *SRegion) GetWires(vpcId, hostId string) ([]SWire, error) {
wires := []SWire{}
params := map[string]interface{}{}
params := map[string]interface{}{"cloud_env": ""}
if len(vpcId) > 0 {
params["vpc_id"] = vpcId
}
+4
View File
@@ -80,3 +80,7 @@ func (self *SVpc) AttachInternetGateway(igwId string) error {
func (self *SVpc) CreateINatGateway(opts *cloudprovider.NatGatewayCreateOptions) (cloudprovider.ICloudNatGateway, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "CreateINatGateway")
}
func (self *SVpc) CreateIWire(opts *cloudprovider.SWireCreateOptions) (cloudprovider.ICloudWire, error) {
return nil, errors.Wrapf(cloudprovider.ErrNotImplemented, "CreateIWire")
}
+1 -1
View File
@@ -210,7 +210,7 @@ func handleServerRemoteConsole(ctx context.Context, w http.ResponseWriter, r *ht
return
}
switch info.Protocol {
case session.ALIYUN, session.QCLOUD, session.OPENSTACK, session.VMRC, session.ZSTACK, session.CTYUN, session.HUAWEI, session.APSARA, session.JDCLOUD:
case session.ALIYUN, session.QCLOUD, session.OPENSTACK, session.VMRC, session.ZSTACK, session.CTYUN, session.HUAWEI, session.APSARA, session.JDCLOUD, session.CLOUDPODS:
responsePublicCloudConsole(ctx, info, w)
case session.VNC, session.SPICE, session.WMKS:
handleDataSession(ctx, info, w, url.Values{"password": {info.GetPassword()}}, true)
+19 -3
View File
@@ -38,6 +38,7 @@ const (
HUAWEI = api.HUAWEI
APSARA = api.APSARA
JDCLOUD = api.JDCLOUD
CLOUDPODS = api.CLOUDPODS
)
type RemoteConsoleInfo struct {
@@ -49,9 +50,12 @@ type RemoteConsoleInfo struct {
VncPassword string `json:"vncPassword"`
// used by aliyun server
InstanceId string `json:"instance_id"`
Url string `json:"url"`
Password string `json:"password"`
InstanceId string `json:"instance_id"`
InstanceName string `json:"instance_name"`
Url string `json:"url"`
Password string `json:"password"`
ConnectParams string `json:"connect_params"`
ApiServer string `json:"api_server"`
}
func NewRemoteConsoleInfoByCloud(s *mcclient.ClientSession, serverId string) (*RemoteConsoleInfo, error) {
@@ -127,6 +131,8 @@ func (info *RemoteConsoleInfo) GetConnectParams() (string, error) {
return info.getApsaraURL()
case QCLOUD:
return info.getQcloudURL()
case CLOUDPODS:
return info.getCloudpodsURL()
case OPENSTACK, VMRC, ZSTACK, CTYUN, HUAWEI, JDCLOUD:
return info.Url, nil
default:
@@ -174,6 +180,16 @@ func (info *RemoteConsoleInfo) getAliyunURL() (string, error) {
return info.getConnParamsURL(base, params), nil
}
func (info *RemoteConsoleInfo) getCloudpodsURL() (string, error) {
base := fmt.Sprintf("%s/web-console/no-vnc", info.ApiServer)
params := url.Values{
"data": {info.ConnectParams},
"instanceId": {info.InstanceId},
"instanceName": {info.InstanceName},
}
return info.getConnParamsURL(base, params), nil
}
func (info *RemoteConsoleInfo) getApsaraURL() (string, error) {
isWindows := "False"
if info.OsName == "Windows" {