Merge pull request #467 in YUNIONIO/onecloud from ~LIZEXI/onecloud:feature/lzx-schedule-disk to release/2.3.0

* commit '0244e70516fa14222350673c757ba279d687eb27':
  support schedule disk
This commit is contained in:
李泽玺
2018-11-12 21:16:39 +08:00
10 changed files with 419 additions and 189 deletions
+10 -3
View File
@@ -166,14 +166,18 @@ func init() {
})
type DiskCreateOptions struct {
STORAGE string `help:"ID or name of storage where the disk is created"`
options.ScheduleOptions
NAME string `help:"Name of the disk"`
DISKDESC string `help:"Image size or size of virtual disk"`
Desc string `help:"Description" metavar:"Description"`
Storage string `help:"ID or name of storage where the disk is created"`
TaskNotify bool `help:"Setup task notify"`
}
R(&DiskCreateOptions{}, "disk-create", "Create a virtual disk", func(s *mcclient.ClientSession, args *DiskCreateOptions) error {
params := jsonutils.NewDict()
params, err := args.ScheduleOptions.Params()
if err != nil {
return err
}
params.Add(jsonutils.NewString(args.NAME), "name")
params.Add(jsonutils.NewString(args.DISKDESC), "disk")
if len(args.Desc) > 0 {
@@ -182,7 +186,10 @@ func init() {
if args.TaskNotify {
s.PrepareTask()
}
disk, err := modules.Disks.CreateInContext(s, params, &modules.Storages, args.STORAGE)
if args.Storage != "" {
params.Add(jsonutils.NewString(args.Storage), "storage")
}
disk, err := modules.Disks.Create(s, params)
if err != nil {
return err
}
+129 -65
View File
@@ -93,7 +93,7 @@ type SDisk struct {
AutoDelete bool `nullable:"false" default:"false" get:"user" update:"user"` // Column(Boolean, nullable=False, default=False)
StorageId string `width:"128" charset:"ascii" nullable:"false" list:"admin" create:"required"` // Column(VARCHAR(ID_LENGTH, charset='ascii'), nullable=False)
StorageId string `width:"128" charset:"ascii" nullable:"true" list:"admin"` // Column(VARCHAR(ID_LENGTH, charset='ascii'), nullable=True)
// # backing template id and type
TemplateId string `width:"256" charset:"ascii" nullable:"true" list:"user"` // Column(VARCHAR(ID_LENGTH, charset='ascii'), nullable=True)
@@ -236,74 +236,136 @@ func (self *SDisk) CustomizeCreate(ctx context.Context, userCred mcclient.TokenC
}
func (manager *SDiskManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerProjId string, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
if disk, err := data.Get("disk"); err != nil {
disk, err := data.Get("disk")
if err != nil {
return nil, err
} else {
if diskConfig, err := parseDiskInfo(ctx, userCred, disk); err != nil {
return nil, err
} else {
data.Add(jsonutils.Marshal(diskConfig), "disk")
if storageID, err := data.GetString("storage_id"); err != nil {
return nil, err
} else {
storages := StorageManager.Query().SubQuery()
storage := SStorage{}
storage.SetModelManager(StorageManager)
if err := storages.Query().Equals("id", storageID).First(&storage); err != nil {
return nil, err
}
if !storage.Enabled {
return nil, httperrors.NewInputParameterError("Cannot create disk with disabled storage[%s]", storage.Name)
}
if !utils.IsInStringArray(storage.Status, []string{STORAGE_ENABLED, STORAGE_ONLINE}) {
return nil, httperrors.NewInputParameterError("Cannot create disk with offline storage[%s]", storage.Name)
}
if len(diskConfig.Backend) == 0 {
diskConfig.Backend = storage.StorageType
}
if storage.StorageType != diskConfig.Backend {
return nil, httperrors.NewInputParameterError("Storage type[%s] not match backend %s", storage.StorageType, diskConfig.Backend)
}
size := diskConfig.Size >> 10
if storage.StorageType == STORAGE_RBD {
diskConfig.Format = "raw"
data.Add(jsonutils.Marshal(diskConfig), "disk")
} else if storage.StorageType == STORAGE_CLOUD_EFFICIENCY || storage.StorageType == STORAGE_CLOUD_SSD {
if size < 20 || size > 32768 {
return nil, httperrors.NewInputParameterError("cloud_ssd or cloud_efficiency disk only support 20G ~ 32768G")
}
} else if storage.StorageType == STORAGE_PUBLIC_CLOUD {
if size < 5 || size > 2000 {
return nil, httperrors.NewInputParameterError("cloud disk only support 5G ~ 2000G")
}
}
hoststorages := HoststorageManager.Query().SubQuery()
hoststorage := make([]SHoststorage, 0)
if err := hoststorages.Query().Equals("storage_id", storage.Id).All(&hoststorage); err != nil {
return nil, err
}
if len(hoststorage) == 0 {
return nil, httperrors.NewInputParameterError("Storage[%s] must attach to a host", storage.Name)
}
if diskConfig.Size > storage.GetFreeCapacity() && !storage.IsEmulated {
return nil, httperrors.NewInputParameterError("Not enough free space")
}
if _, err := manager.SSharableVirtualResourceBaseManager.ValidateCreateData(ctx, userCred, ownerProjId, query, data); err != nil {
return nil, err
}
pendingUsage := SQuota{Storage: diskConfig.Size}
if err := QuotaManager.CheckSetPendingQuota(ctx, userCred, userCred.GetProjectId(), &pendingUsage); err != nil {
return nil, err
}
}
}
diskConfig, err := parseDiskInfo(ctx, userCred, disk)
if err != nil {
return nil, err
}
storageID := jsonutils.GetAnyString(data, []string{"storage_id", "storage"})
if storageID != "" {
storageObj, err := StorageManager.FetchByIdOrName(nil, storageID)
if err != nil {
return nil, httperrors.NewResourceNotFoundError("Storage %s not found", storageID)
}
storage := storageObj.(*SStorage)
if len(diskConfig.Backend) == 0 {
diskConfig.Backend = storage.StorageType
}
err = manager.validateDiskOnStorage(diskConfig, storage)
if err != nil {
return nil, err
}
data.Add(jsonutils.NewString(storage.Id), "storage_id")
} else {
diskConfig.Backend = STORAGE_LOCAL
hypervisor, _ := data.GetString("hypervisor")
data, err = ValidateScheduleCreateData(ctx, userCred, data, hypervisor)
if err != nil {
return nil, err
}
}
data.Add(jsonutils.Marshal(diskConfig), "disk")
if _, err := manager.SSharableVirtualResourceBaseManager.ValidateCreateData(ctx, userCred, ownerProjId, query, data); err != nil {
return nil, err
}
pendingUsage := SQuota{Storage: diskConfig.Size}
if err := QuotaManager.CheckSetPendingQuota(ctx, userCred, userCred.GetProjectId(), &pendingUsage); err != nil {
return nil, err
}
return data, nil
}
func (disk *SDisk) PostCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerProjId string, query jsonutils.JSONObject, data jsonutils.JSONObject) {
disk.SSharableVirtualResourceBase.PostCreate(ctx, userCred, ownerProjId, query, data)
disk.StartDiskCreateTask(ctx, userCred, false, "", "")
func (manager *SDiskManager) validateDiskOnStorage(diskConfig *SDiskConfig, storage *SStorage) error {
if !storage.Enabled {
return httperrors.NewInputParameterError("Cannot create disk with disabled storage[%s]", storage.Name)
}
if !utils.IsInStringArray(storage.Status, []string{STORAGE_ENABLED, STORAGE_ONLINE}) {
return httperrors.NewInputParameterError("Cannot create disk with offline storage[%s]", storage.Name)
}
if storage.StorageType != diskConfig.Backend {
return httperrors.NewInputParameterError("Storage type[%s] not match backend %s", storage.StorageType, diskConfig.Backend)
}
size := diskConfig.Size >> 10
if storage.StorageType == STORAGE_CLOUD_EFFICIENCY || storage.StorageType == STORAGE_CLOUD_SSD {
if size < 20 || size > 32768 {
return httperrors.NewInputParameterError("cloud_ssd or cloud_efficiency disk only support 20G ~ 32768G")
}
} else if storage.StorageType == STORAGE_PUBLIC_CLOUD {
if size < 5 || size > 2000 {
return httperrors.NewInputParameterError("cloud disk only support 5G ~ 2000G")
}
}
hoststorages := HoststorageManager.Query().SubQuery()
hoststorage := make([]SHoststorage, 0)
if err := hoststorages.Query().Equals("storage_id", storage.Id).All(&hoststorage); err != nil {
return err
}
if len(hoststorage) == 0 {
return httperrors.NewInputParameterError("Storage[%s] must attach to a host", storage.Name)
}
if diskConfig.Size > storage.GetFreeCapacity() && !storage.IsEmulated {
return httperrors.NewInputParameterError("Not enough free space")
}
return nil
}
func (disk *SDisk) SetStorageByHost(hostId string, diskConfig *SDiskConfig) error {
host := HostManager.FetchHostById(hostId)
backend := diskConfig.Backend
if backend == "" {
return fmt.Errorf("Backend is empty")
}
var storage *SStorage
if utils.IsInStringArray(backend, STORAGE_LIMITED_TYPES) {
storage = host.GetLeastUsedStorage(backend)
} else {
// unlimited pulic cloud storages
storages := host.GetAttachedStorages("")
for _, s := range storages {
if s.StorageType == backend {
tmpS := s
storage = &tmpS
}
}
}
if storage == nil {
return fmt.Errorf("Not found host %s backend %s storage", host.Name, backend)
}
err := DiskManager.validateDiskOnStorage(diskConfig, storage)
if err != nil {
return err
}
_, err = disk.GetModelManager().TableSpec().Update(disk, func() error {
disk.StorageId = storage.Id
return nil
})
return err
}
func getDiskResourceRequirements(ctx context.Context, userCred mcclient.TokenCredential, data jsonutils.JSONObject, count int) SQuota {
diskSize, _ := data.Int("disk", "size")
return SQuota{
Storage: int(diskSize) * count,
}
}
func (manager *SDiskManager) convertToBatchCreateData(data jsonutils.JSONObject) *jsonutils.JSONDict {
diskConfig, _ := data.Get("disk")
newData := data.(*jsonutils.JSONDict).CopyExcludes("disk")
newData.Add(diskConfig, "disk.0")
return newData
}
func (manager *SDiskManager) OnCreateComplete(ctx context.Context, items []db.IModel, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) {
pendingUsage := getDiskResourceRequirements(ctx, userCred, data, len(items))
RunBatchCreateTask(ctx, items, userCred, manager.convertToBatchCreateData(data), pendingUsage, "DiskBatchCreateTask")
}
func (self *SDisk) StartDiskCreateTask(ctx context.Context, userCred mcclient.TokenCredential, rebuild bool, snapshot string, parentTaskId string) error {
@@ -1176,8 +1238,10 @@ func (self *SDisk) GetShortDesc() *jsonutils.JSONDict {
desc := self.SSharableVirtualResourceBase.GetShortDesc()
desc.Add(jsonutils.NewInt(int64(self.DiskSize)), "size")
storage := self.GetStorage()
desc.Add(jsonutils.NewString(storage.StorageType), "storage_type")
desc.Add(jsonutils.NewString(storage.MediumType), "medium_type")
if storage != nil {
desc.Add(jsonutils.NewString(storage.StorageType), "storage_type")
desc.Add(jsonutils.NewString(storage.MediumType), "medium_type")
}
if priceKey := self.GetMetadata("price_key", nil); len(priceKey) > 0 {
desc.Add(jsonutils.NewString(priceKey), "price_key")
+5 -117
View File
@@ -661,112 +661,11 @@ func (manager *SGuestManager) ValidateCreateData(ctx context.Context, userCred m
}
}
if jsonutils.QueryBoolean(data, "baremetal", false) {
hypervisor = HYPERVISOR_BAREMETAL
data, err = ValidateScheduleCreateData(ctx, userCred, data, hypervisor)
if err != nil {
return nil, err
}
// base validate_create_data
if (data.Contains("prefer_baremetal") || data.Contains("prefer_host")) && hypervisor != HYPERVISOR_CONTAINER {
if !userCred.IsSystemAdmin() {
return nil, httperrors.NewNotSufficientPrivilegeError("Only system admin can specify preferred host")
}
bmName, _ := data.GetString("prefer_host")
if len(bmName) == 0 {
bmName, _ = data.GetString("prefer_baremetal")
}
bmObj, err := HostManager.FetchByIdOrName(nil, bmName)
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError("Host %s not found", bmName)
} else {
return nil, httperrors.NewGeneralError(err)
}
}
baremetal := bmObj.(*SHost)
if !baremetal.Enabled {
return nil, httperrors.NewInvalidStatusError("Baremetal %s not enabled", bmName)
}
if len(hypervisor) > 0 && hypervisor != HOSTTYPE_HYPERVISOR[baremetal.HostType] {
return nil, httperrors.NewInputParameterError("cannot run hypervisor %s on specified host with type %s", hypervisor, baremetal.HostType)
}
if len(hypervisor) == 0 {
hypervisor = HOSTTYPE_HYPERVISOR[baremetal.HostType]
}
if len(hypervisor) == 0 {
hypervisor = HYPERVISOR_DEFAULT
}
_, err = GetDriver(hypervisor).ValidateCreateHostData(ctx, userCred, bmName, baremetal, data)
if err != nil {
return nil, err
}
} else {
schedtags := make(map[string]string)
if data.Contains("aggregate_strategy") {
err = data.Unmarshal(&schedtags, "aggregate_strategy")
if err != nil {
return nil, httperrors.NewInputParameterError("invalid aggregate_strategy")
}
}
for idx := 0; data.Contains(fmt.Sprintf("schedtag.%d", idx)); idx += 1 {
aggStr, _ := data.GetString(fmt.Sprintf("schedtag.%d", idx))
if len(aggStr) > 0 {
parts := strings.Split(aggStr, ":")
if len(parts) >= 2 && len(parts[0]) > 0 && len(parts[1]) > 0 {
schedtags[parts[0]] = parts[1]
}
}
}
if len(schedtags) > 0 {
schedtags, err = SchedtagManager.ValidateSchedtags(userCred, schedtags)
if err != nil {
return nil, httperrors.NewInputParameterError("invalid aggregate_strategy: %s", err)
}
data.Add(jsonutils.Marshal(schedtags), "aggregate_strategy")
}
if data.Contains("prefer_wire") {
wireStr, _ := data.GetString("prefer_wire")
wireObj, err := WireManager.FetchById(wireStr)
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError("Wire %s not found", wireStr)
} else {
return nil, httperrors.NewGeneralError(err)
}
}
wire := wireObj.(*SWire)
data.Add(jsonutils.NewString(wire.Id), "prefer_wire_id")
zone := wire.GetZone()
data.Add(jsonutils.NewString(zone.Id), "prefer_zone_id")
} else if data.Contains("prefer_zone") {
zoneStr, _ := data.GetString("prefer_zone")
zoneObj, err := ZoneManager.FetchById(zoneStr)
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError("Zone %s not found", zoneStr)
} else {
return nil, httperrors.NewGeneralError(err)
}
}
zone := zoneObj.(*SZone)
data.Add(jsonutils.NewString(zone.Id), "prefer_zone_id")
}
}
// default hypervisor
if len(hypervisor) == 0 {
hypervisor = HYPERVISOR_KVM
}
if !utils.IsInStringArray(hypervisor, HYPERVISORS) {
return nil, httperrors.NewInputParameterError("Hypervisor %s not supported", hypervisor)
}
data.Add(jsonutils.NewString(hypervisor), "hypervisor")
hypervisor, _ = data.GetString("hypervisor")
for idx := 0; data.Contains(fmt.Sprintf("net.%d", idx)); idx += 1 {
netJson, err := data.Get(fmt.Sprintf("net.%d", idx))
@@ -986,18 +885,7 @@ func (guest *SGuest) setApptags(ctx context.Context, appTags []string, userCred
func (manager *SGuestManager) OnCreateComplete(ctx context.Context, items []db.IModel, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) {
pendingUsage := getGuestResourceRequirements(ctx, userCred, data, len(items))
taskItems := make([]db.IStandaloneModel, len(items))
for i, t := range items {
taskItems[i] = t.(db.IStandaloneModel)
}
params := data.(*jsonutils.JSONDict)
task, err := taskman.TaskManager.NewParallelTask(ctx, "GuestBatchCreateTask", taskItems, userCred, params, "", "", &pendingUsage)
if err != nil {
log.Errorf("GuestBatchCreateTask newTask error %s", err)
} else {
task.ScheduleRun(nil)
}
RunBatchCreateTask(ctx, items, userCred, data, pendingUsage, "GuestBatchCreateTask")
}
func (guest *SGuest) GetGroups() []SGroupguest {
+151
View File
@@ -0,0 +1,151 @@
package models
import (
"context"
"database/sql"
"fmt"
"strings"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/utils"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
)
func RunBatchCreateTask(
ctx context.Context,
items []db.IModel,
userCred mcclient.TokenCredential,
data jsonutils.JSONObject,
pendingUsage SQuota,
taskName string,
) {
taskItems := make([]db.IStandaloneModel, len(items))
for i, t := range items {
taskItems[i] = t.(db.IStandaloneModel)
}
params := data.(*jsonutils.JSONDict)
task, err := taskman.TaskManager.NewParallelTask(ctx, taskName, taskItems, userCred, params, "", "", &pendingUsage)
if err != nil {
log.Errorf("%s newTask error %s", taskName, err)
} else {
task.ScheduleRun(nil)
}
}
func ValidateScheduleCreateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict, hypervisor string) (*jsonutils.JSONDict, error) {
var err error
if jsonutils.QueryBoolean(data, "baremetal", false) {
hypervisor = HYPERVISOR_BAREMETAL
}
// base validate_create_data
if (data.Contains("prefer_baremetal") || data.Contains("prefer_host")) && hypervisor != HYPERVISOR_CONTAINER {
if !userCred.IsSystemAdmin() {
return nil, httperrors.NewNotSufficientPrivilegeError("Only system admin can specify preferred host")
}
bmName, _ := data.GetString("prefer_host")
if len(bmName) == 0 {
bmName, _ = data.GetString("prefer_baremetal")
}
bmObj, err := HostManager.FetchByIdOrName(nil, bmName)
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError("Host %s not found", bmName)
} else {
return nil, httperrors.NewGeneralError(err)
}
}
baremetal := bmObj.(*SHost)
if !baremetal.Enabled {
return nil, httperrors.NewInvalidStatusError("Baremetal %s not enabled", bmName)
}
if len(hypervisor) > 0 && hypervisor != HOSTTYPE_HYPERVISOR[baremetal.HostType] {
return nil, httperrors.NewInputParameterError("cannot run hypervisor %s on specified host with type %s", hypervisor, baremetal.HostType)
}
if len(hypervisor) == 0 {
hypervisor = HOSTTYPE_HYPERVISOR[baremetal.HostType]
}
if len(hypervisor) == 0 {
hypervisor = HYPERVISOR_DEFAULT
}
_, err = GetDriver(hypervisor).ValidateCreateHostData(ctx, userCred, bmName, baremetal, data)
if err != nil {
return nil, err
}
} else {
schedtags := make(map[string]string)
if data.Contains("aggregate_strategy") {
err = data.Unmarshal(&schedtags, "aggregate_strategy")
if err != nil {
return nil, httperrors.NewInputParameterError("invalid aggregate_strategy")
}
}
for idx := 0; data.Contains(fmt.Sprintf("schedtag.%d", idx)); idx += 1 {
aggStr, _ := data.GetString(fmt.Sprintf("schedtag.%d", idx))
if len(aggStr) > 0 {
parts := strings.Split(aggStr, ":")
if len(parts) >= 2 && len(parts[0]) > 0 && len(parts[1]) > 0 {
schedtags[parts[0]] = parts[1]
}
}
}
if len(schedtags) > 0 {
schedtags, err = SchedtagManager.ValidateSchedtags(userCred, schedtags)
if err != nil {
return nil, httperrors.NewInputParameterError("invalid aggregate_strategy: %s", err)
}
data.Add(jsonutils.Marshal(schedtags), "aggregate_strategy")
}
if data.Contains("prefer_wire") {
wireStr, _ := data.GetString("prefer_wire")
wireObj, err := WireManager.FetchById(wireStr)
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError("Wire %s not found", wireStr)
} else {
return nil, httperrors.NewGeneralError(err)
}
}
wire := wireObj.(*SWire)
data.Add(jsonutils.NewString(wire.Id), "prefer_wire_id")
zone := wire.GetZone()
data.Add(jsonutils.NewString(zone.Id), "prefer_zone_id")
} else if data.Contains("prefer_zone") {
zoneStr, _ := data.GetString("prefer_zone")
zoneObj, err := ZoneManager.FetchById(zoneStr)
if err != nil {
if err == sql.ErrNoRows {
return nil, httperrors.NewResourceNotFoundError("Zone %s not found", zoneStr)
} else {
return nil, httperrors.NewGeneralError(err)
}
}
zone := zoneObj.(*SZone)
data.Add(jsonutils.NewString(zone.Id), "prefer_zone_id")
}
}
// default hypervisor
if len(hypervisor) == 0 {
hypervisor = HYPERVISOR_KVM
}
if !utils.IsInStringArray(hypervisor, HYPERVISORS) {
return nil, httperrors.NewInputParameterError("Hypervisor %s not supported", hypervisor)
}
data.Add(jsonutils.NewString(hypervisor), "hypervisor")
return data, nil
}
+10
View File
@@ -149,3 +149,13 @@ func (self *SHoststorage) Delete(ctx context.Context, userCred mcclient.TokenCre
func (self *SHoststorage) Detach(ctx context.Context, userCred mcclient.TokenCredential) error {
return db.DetachJoint(ctx, userCred, self)
}
func (manager *SHoststorageManager) GetStorages(hostId string) ([]SHoststorage, error) {
hoststorage := make([]SHoststorage, 0)
hoststorages := HoststorageManager.Query().SubQuery()
err := hoststorages.Query().Equals("host_id", hostId).All(&hoststorage)
if err != nil {
return nil, err
}
return hoststorage, nil
}
+1
View File
@@ -53,6 +53,7 @@ var (
STORAGE_LOCAL, STORAGE_BAREMETAL, STORAGE_SHEEPDOG,
STORAGE_RBD, STORAGE_DOCKER, STORAGE_NAS, STORAGE_VSAN,
}
STORAGE_LIMITED_TYPES = []string{STORAGE_LOCAL, STORAGE_BAREMETAL, STORAGE_NAS, STORAGE_RBD}
)
type SStorageManager struct {
@@ -0,0 +1,91 @@
package tasks
import (
"context"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudcommon/notifyclient"
"yunion.io/x/onecloud/pkg/compute/models"
)
type DiskBatchCreateTask struct {
taskman.STask
}
func init() {
taskman.RegisterTask(DiskBatchCreateTask{})
}
func (self *DiskBatchCreateTask) getNeedScheduleDisks(objs []db.IStandaloneModel) []db.IStandaloneModel {
toSchedDisks := make([]db.IStandaloneModel, 0)
for _, obj := range objs {
disk := obj.(*models.SDisk)
if disk.StorageId == "" {
toSchedDisks = append(toSchedDisks, disk)
}
}
return toSchedDisks
}
func (self *DiskBatchCreateTask) OnInit(ctx context.Context, objs []db.IStandaloneModel, body jsonutils.JSONObject) {
toSchedDisks := self.getNeedScheduleDisks(objs)
if len(toSchedDisks) == 0 {
self.SetStage("OnScheduleComplete", nil)
// create not need schedule disks directly
for _, disk := range objs {
self.startCreateDisk(ctx, disk.(*models.SDisk))
}
return
}
StartScheduleObjects(ctx, self, toSchedDisks)
}
func (self *DiskBatchCreateTask) OnScheduleFailCallback(obj IScheduleModel) {
disk := obj.(*models.SDisk)
log.Errorf("Schedule disk %s failed", disk.Name)
}
func (self *DiskBatchCreateTask) SaveScheduleResult(ctx context.Context, obj IScheduleModel, hostId string) {
var err error
disk := obj.(*models.SDisk)
pendingUsage := models.SQuota{}
err = self.GetPendingUsage(&pendingUsage)
if err != nil {
log.Errorf("GetPendingUsage fail %s", err)
}
diskConfig := models.SDiskConfig{}
self.GetParams().Unmarshal(&diskConfig, "disk.0")
quotaStorage := models.SQuota{Storage: disk.DiskSize}
err = disk.SetStorageByHost(hostId, &diskConfig)
if err != nil {
models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, disk.ProjectId, &pendingUsage, &quotaStorage)
disk.SetStatus(self.UserCred, models.DISK_ALLOC_FAILED, err.Error())
self.SetStageFailed(ctx, err.Error())
db.OpsLog.LogEvent(disk, db.ACT_ALLOCATE_FAIL, err, self.UserCred)
notifyclient.NotifySystemError(disk.Id, disk.Name, models.DISK_ALLOC_FAILED, err.Error())
return
}
self.startCreateDisk(ctx, disk)
}
func (self *DiskBatchCreateTask) startCreateDisk(ctx context.Context, disk *models.SDisk) {
pendingUsage := models.SQuota{}
err := self.GetPendingUsage(&pendingUsage)
if err != nil {
log.Errorf("GetPendingUsage fail %s", err)
}
quotaStorage := models.SQuota{Storage: disk.DiskSize}
models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, disk.ProjectId, &pendingUsage, &quotaStorage)
self.SetPendingUsage(&pendingUsage)
disk.StartDiskCreateTask(ctx, self.GetUserCred(), false, "", self.GetTaskId())
}
func (self *DiskBatchCreateTask) OnScheduleComplete(ctx context.Context, items []db.IStandaloneModel, data *jsonutils.JSONDict) {
self.SetStageComplete(ctx, nil)
}
@@ -5,6 +5,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudcommon/notifyclient"
+14
View File
@@ -0,0 +1,14 @@
package options
import "yunion.io/x/jsonutils"
type ScheduleOptions struct {
Zone string `help:"Preferred zone where virtual server should be created" json:"prefer_zone"`
Host string `help:"Preferred host where virtual server should be created" json:"prefer_host"`
Schedtag []string `help:"Schedule policy, key = aggregate name, value = require|exclude|prefer|avoid" metavar:"<KEY:VALUE>"`
Hypervisor string `help:"Hypervisor type" choices:"kvm|esxi|baremetal|container|aliyun|azure|qcloud"`
}
func (opts *ScheduleOptions) Params() (*jsonutils.JSONDict, error) {
return optionsStructToParams(opts)
}
+7 -4
View File
@@ -89,6 +89,7 @@ func ParseServerDeployInfoList(list []string) (*jsonutils.JSONDict, error) {
}
type ServerCreateOptions struct {
ScheduleOptions
NAME string `help:"Name of server"`
MEM string `help:"Memory size" metavar:"MEMORY" json:"vmem_size"`
Disk []string `help:"Disk descriptions" nargs:"+"`
@@ -107,15 +108,11 @@ type ServerCreateOptions struct {
AllowDelete *bool `help:"Unlock server to allow deleting" json:"-"`
ShutdownBehavior string `help:"Behavior after VM server shutdown, stop or terminate server" metavar:"<SHUTDOWN_BEHAVIOR>" choices:"stop|terminate"`
AutoStart *bool `help:"Auto start server after it is created"`
Zone string `help:"Preferred zone where virtual server should be created" json:"prefer_zone"`
Host string `help:"Preferred host where virtual server should be created" json:"prefer_host"`
Schedtag []string `help:"Schedule policy, key = aggregate name, value = require|exclude|prefer|avoid" metavar:"<KEY:VALUE>"`
Deploy []string `help:"Specify deploy files in virtual server file system" json:"-"`
Group []string `help:"Group of virtual server"`
Project string `help:"'Owner project ID or Name" json:"tenant"`
User string `help:"Owner user ID or Name"`
System *bool `help:"Create a system VM, sysadmin ONLY option" json:"is_system"`
Hypervisor string `help:"Hypervisor type" choices:"kvm|esxi|baremetal|container|aliyun|azure|qcloud"`
TaskNotify *bool `help:"Setup task notify" json:"-"`
Count *int `help:"Create multiple simultaneously" default:"1" json:"-"`
DryRun *bool `help:"Dry run to test scheduler" json:"-"`
@@ -124,10 +121,16 @@ type ServerCreateOptions struct {
}
func (opts *ServerCreateOptions) Params() (*jsonutils.JSONDict, error) {
schedParams, err := opts.ScheduleOptions.Params()
if err != nil {
return nil, err
}
params, err := optionsStructToParams(opts)
if err != nil {
return nil, err
}
params.Update(schedParams)
{
deployParams, err := ParseServerDeployInfoList(opts.Deploy)