diff --git a/cmd/climc/shell/disks.go b/cmd/climc/shell/disks.go index 0e559cdbaf..d152e971f7 100644 --- a/cmd/climc/shell/disks.go +++ b/cmd/climc/shell/disks.go @@ -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 } diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 73c5677b6f..6c11a9f2a3 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -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") diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 7a1e4c4de0..50b823721c 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -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 { diff --git a/pkg/compute/models/helper.go b/pkg/compute/models/helper.go new file mode 100644 index 0000000000..41e680cdf6 --- /dev/null +++ b/pkg/compute/models/helper.go @@ -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 +} diff --git a/pkg/compute/models/hoststorages.go b/pkg/compute/models/hoststorages.go index 16bf9e7fee..8da2193ddb 100644 --- a/pkg/compute/models/hoststorages.go +++ b/pkg/compute/models/hoststorages.go @@ -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 +} diff --git a/pkg/compute/models/storages.go b/pkg/compute/models/storages.go index c8124b2956..b7be546d7d 100644 --- a/pkg/compute/models/storages.go +++ b/pkg/compute/models/storages.go @@ -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 { diff --git a/pkg/compute/tasks/disk_batch_create_task.go b/pkg/compute/tasks/disk_batch_create_task.go new file mode 100644 index 0000000000..26b4d273fd --- /dev/null +++ b/pkg/compute/tasks/disk_batch_create_task.go @@ -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, "aStorage) + 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, "aStorage) + 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) +} diff --git a/pkg/compute/tasks/guest_batch_create_task.go b/pkg/compute/tasks/guest_batch_create_task.go index dab18475c3..8df1e5cbab 100644 --- a/pkg/compute/tasks/guest_batch_create_task.go +++ b/pkg/compute/tasks/guest_batch_create_task.go @@ -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" diff --git a/pkg/mcclient/options/schedule.go b/pkg/mcclient/options/schedule.go new file mode 100644 index 0000000000..d76d9b923c --- /dev/null +++ b/pkg/mcclient/options/schedule.go @@ -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:""` + Hypervisor string `help:"Hypervisor type" choices:"kvm|esxi|baremetal|container|aliyun|azure|qcloud"` +} + +func (opts *ScheduleOptions) Params() (*jsonutils.JSONDict, error) { + return optionsStructToParams(opts) +} diff --git a/pkg/mcclient/options/servers.go b/pkg/mcclient/options/servers.go index 69793506e6..89924038bf 100644 --- a/pkg/mcclient/options/servers.go +++ b/pkg/mcclient/options/servers.go @@ -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:"" 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:""` 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)