diff --git a/pkg/apis/compute/api.go b/pkg/apis/compute/api.go index ab775eee48..2bb893a810 100644 --- a/pkg/apis/compute/api.go +++ b/pkg/apis/compute/api.go @@ -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 diff --git a/pkg/apis/compute/guests.go b/pkg/apis/compute/guests.go index dc5b6be64c..22cced81ae 100644 --- a/pkg/apis/compute/guests.go +++ b/pkg/apis/compute/guests.go @@ -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"` +} diff --git a/pkg/apis/webconsole/consts.go b/pkg/apis/webconsole/consts.go index c4ddbcc0d8..c484060fc5 100644 --- a/pkg/apis/webconsole/consts.go +++ b/pkg/apis/webconsole/consts.go @@ -32,4 +32,5 @@ const ( HUAWEI = "huawei" APSARA = "apsara" JDCLOUD = "jdcloud" + CLOUDPODS = "cloudpods" ) diff --git a/pkg/cloudprovider/network.go b/pkg/cloudprovider/network.go index 3e863d6495..f323590f2e 100644 --- a/pkg/cloudprovider/network.go +++ b/pkg/cloudprovider/network.go @@ -20,3 +20,10 @@ type SNetworkCreateOptions struct { ProjectId string Cidr string } + +type SWireCreateOptions struct { + Name string + ZoneId string + Bandwidth int + Mtu int +} diff --git a/pkg/cloudprovider/resources.go b/pkg/cloudprovider/resources.go index f24459aebb..7a493fc561 100644 --- a/pkg/cloudprovider/resources.go +++ b/pkg/cloudprovider/resources.go @@ -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) diff --git a/pkg/compute/guestdrivers/azure.go b/pkg/compute/guestdrivers/azure.go index bdcdde2bbb..a92bbe61c7 100644 --- a/pkg/compute/guestdrivers/azure.go +++ b/pkg/compute/guestdrivers/azure.go @@ -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 } diff --git a/pkg/compute/guestdrivers/baremetals.go b/pkg/compute/guestdrivers/baremetals.go index 773114a677..43c7f04e7d 100644 --- a/pkg/compute/guestdrivers/baremetals.go +++ b/pkg/compute/guestdrivers/baremetals.go @@ -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 { diff --git a/pkg/compute/guestdrivers/base.go b/pkg/compute/guestdrivers/base.go index 356aeeb33b..40f3a8be00 100644 --- a/pkg/compute/guestdrivers/base.go +++ b/pkg/compute/guestdrivers/base.go @@ -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") } diff --git a/pkg/compute/guestdrivers/cloudpods.go b/pkg/compute/guestdrivers/cloudpods.go index d4d041dbe4..a445cf1907 100644 --- a/pkg/compute/guestdrivers/cloudpods.go +++ b/pkg/compute/guestdrivers/cloudpods.go @@ -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) diff --git a/pkg/compute/guestdrivers/container.go b/pkg/compute/guestdrivers/container.go index f0fbe99154..2cfeac276a 100644 --- a/pkg/compute/guestdrivers/container.go +++ b/pkg/compute/guestdrivers/container.go @@ -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 { diff --git a/pkg/compute/guestdrivers/google.go b/pkg/compute/guestdrivers/google.go index 2f58023053..9f047c6a51 100644 --- a/pkg/compute/guestdrivers/google.go +++ b/pkg/compute/guestdrivers/google.go @@ -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) { diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index 376aa1b026..1208e42612 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -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) { diff --git a/pkg/compute/guestdrivers/managedvirtual.go b/pkg/compute/guestdrivers/managedvirtual.go index ccbf3a2906..97d739bfbd 100644 --- a/pkg/compute/guestdrivers/managedvirtual.go +++ b/pkg/compute/guestdrivers/managedvirtual.go @@ -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 diff --git a/pkg/compute/guestdrivers/virtualization.go b/pkg/compute/guestdrivers/virtualization.go index ffe3aeccfe..d026e4690a 100644 --- a/pkg/compute/guestdrivers/virtualization.go +++ b/pkg/compute/guestdrivers/virtualization.go @@ -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 { diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 215fdeceab..67cddb76aa 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -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 } diff --git a/pkg/compute/models/guestdrivers.go b/pkg/compute/models/guestdrivers.go index de42383aed..302910a305 100644 --- a/pkg/compute/models/guestdrivers.go +++ b/pkg/compute/models/guestdrivers.go @@ -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 diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 431edf4809..906fae0717 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -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("a, "pending_usage") - quotas.CancelPendingUsage(ctx, userCred, "a, "a, 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) } diff --git a/pkg/compute/models/networks.go b/pkg/compute/models/networks.go index dec63421cf..ec4e434c6b 100644 --- a/pkg/compute/models/networks.go +++ b/pkg/compute/models/networks.go @@ -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) { diff --git a/pkg/compute/models/vpcs.go b/pkg/compute/models/vpcs.go index 65ec53ae3f..8b701cd1d7 100644 --- a/pkg/compute/models/vpcs.go +++ b/pkg/compute/models/vpcs.go @@ -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 diff --git a/pkg/compute/tasks/guest_backup_tasks.go b/pkg/compute/tasks/guest_backup_tasks.go index 486b4b6cbd..9034b6d280 100644 --- a/pkg/compute/tasks/guest_backup_tasks.go +++ b/pkg/compute/tasks/guest_backup_tasks.go @@ -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())) } } diff --git a/pkg/compute/tasks/guest_start_task.go b/pkg/compute/tasks/guest_start_task.go index 0882b5124e..035d2de3af 100644 --- a/pkg/compute/tasks/guest_start_task.go +++ b/pkg/compute/tasks/guest_start_task.go @@ -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) { diff --git a/pkg/compute/tasks/guest_stop_task.go b/pkg/compute/tasks/guest_stop_task.go index 109d2198fd..a468dea0a3 100644 --- a/pkg/compute/tasks/guest_stop_task.go +++ b/pkg/compute/tasks/guest_stop_task.go @@ -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) { diff --git a/pkg/compute/tasks/ha_guest_start_task.go b/pkg/compute/tasks/ha_guest_start_task.go index 311cc99493..2142224460 100644 --- a/pkg/compute/tasks/ha_guest_start_task.go +++ b/pkg/compute/tasks/ha_guest_start_task.go @@ -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 } } diff --git a/pkg/mcclient/options/servers.go b/pkg/mcclient/options/servers.go index dbe1d02efd..9b75a6512d 100644 --- a/pkg/mcclient/options/servers.go +++ b/pkg/mcclient/options/servers.go @@ -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 diff --git a/pkg/multicloud/cloudpods/cloudpods.go b/pkg/multicloud/cloudpods/cloudpods.go index 0243589e01..1a94403ba6 100644 --- a/pkg/multicloud/cloudpods/cloudpods.go +++ b/pkg/multicloud/cloudpods/cloudpods.go @@ -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 { diff --git a/pkg/multicloud/cloudpods/disk.go b/pkg/multicloud/cloudpods/disk.go index 24b2e0f09a..b8410c99ea 100644 --- a/pkg/multicloud/cloudpods/disk.go +++ b/pkg/multicloud/cloudpods/disk.go @@ -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 diff --git a/pkg/multicloud/cloudpods/eip.go b/pkg/multicloud/cloudpods/eip.go index 59e11bba6b..1c071046e5 100644 --- a/pkg/multicloud/cloudpods/eip.go +++ b/pkg/multicloud/cloudpods/eip.go @@ -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 { diff --git a/pkg/multicloud/cloudpods/instance.go b/pkg/multicloud/cloudpods/instance.go index bdea08688a..355b6e22e3 100644 --- a/pkg/multicloud/cloudpods/instance.go +++ b/pkg/multicloud/cloudpods/instance.go @@ -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, diff --git a/pkg/multicloud/cloudpods/keypair.go b/pkg/multicloud/cloudpods/keypair.go new file mode 100644 index 0000000000..c45bb9407b --- /dev/null +++ b/pkg/multicloud/cloudpods/keypair.go @@ -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 +} diff --git a/pkg/multicloud/cloudpods/region.go b/pkg/multicloud/cloudpods/region.go index 6181625fef..4220961bbf 100644 --- a/pkg/multicloud/cloudpods/region.go +++ b/pkg/multicloud/cloudpods/region.go @@ -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) } diff --git a/pkg/multicloud/cloudpods/storage.go b/pkg/multicloud/cloudpods/storage.go index 75474ee13f..cccbb1e3dc 100644 --- a/pkg/multicloud/cloudpods/storage.go +++ b/pkg/multicloud/cloudpods/storage.go @@ -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 } diff --git a/pkg/multicloud/cloudpods/vpc.go b/pkg/multicloud/cloudpods/vpc.go index 5cc8a511d8..afdd3b46ae 100644 --- a/pkg/multicloud/cloudpods/vpc.go +++ b/pkg/multicloud/cloudpods/vpc.go @@ -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 diff --git a/pkg/multicloud/cloudpods/wire.go b/pkg/multicloud/cloudpods/wire.go index 87a26031ee..c3bff7b63a 100644 --- a/pkg/multicloud/cloudpods/wire.go +++ b/pkg/multicloud/cloudpods/wire.go @@ -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 } diff --git a/pkg/multicloud/vpc_base.go b/pkg/multicloud/vpc_base.go index 4e241b61b4..ab903ed982 100644 --- a/pkg/multicloud/vpc_base.go +++ b/pkg/multicloud/vpc_base.go @@ -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") +} diff --git a/pkg/webconsole/handlers.go b/pkg/webconsole/handlers.go index 2d26177ebc..599061703e 100644 --- a/pkg/webconsole/handlers.go +++ b/pkg/webconsole/handlers.go @@ -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) diff --git a/pkg/webconsole/session/remote_console.go b/pkg/webconsole/session/remote_console.go index 603ec4eeb9..e9cecf40e1 100644 --- a/pkg/webconsole/session/remote_console.go +++ b/pkg/webconsole/session/remote_console.go @@ -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" {