From bc7f62a2add490be9a42385088d4829803ec70af Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 3 Jan 2020 00:40:00 +0800 Subject: [PATCH] scheduler schedule according to quota limit --- pkg/apis/scheduler/api.go | 2 + pkg/cloudcommon/db/quotas/interface.go | 2 + pkg/cloudcommon/db/quotas/quotas.go | 46 ++++++ pkg/cloudcommon/db/quotas/register.go | 5 + pkg/compute/models/cloudaccounts.go | 4 +- pkg/compute/models/guest_actions.go | 2 +- pkg/compute/models/guests.go | 2 +- pkg/compute/models/managedresource.go | 4 +- pkg/compute/models/networks.go | 2 +- pkg/compute/models/projectquota.go | 9 ++ pkg/compute/models/quotas.go | 26 +++- pkg/compute/models/regionquota.go | 39 ++++++ pkg/compute/models/zonequota.go | 4 + pkg/compute/tasks/disk_batch_create_task.go | 21 +++ pkg/compute/tasks/guest_batch_create_task.go | 30 ++++ pkg/compute/tasks/pending_usage.go | 16 ++- pkg/compute/tasks/schedule.go | 49 +++---- pkg/image/models/quotas.go | 9 ++ .../algorithm/predicates/quota_predicate.go | 131 ++++++++++++++++++ pkg/scheduler/algorithmprovider/defaults.go | 1 + pkg/scheduler/cache/candidate/base.go | 29 ++++ pkg/scheduler/core/types.go | 2 + 22 files changed, 392 insertions(+), 43 deletions(-) create mode 100644 pkg/scheduler/algorithm/predicates/quota_predicate.go diff --git a/pkg/apis/scheduler/api.go b/pkg/apis/scheduler/api.go index 80343a62bb..f88336ec28 100644 --- a/pkg/apis/scheduler/api.go +++ b/pkg/apis/scheduler/api.go @@ -77,6 +77,8 @@ type ScheduleInput struct { LiveMigrate bool `json:"live_migrate"` CpuDesc string `json:"cpu_desc"` CpuMicrocode string `json:"cpu_microcode"` + + PendingUsages []jsonutils.JSONObject } func (input ScheduleInput) ToConditionInput() *jsonutils.JSONDict { diff --git a/pkg/cloudcommon/db/quotas/interface.go b/pkg/cloudcommon/db/quotas/interface.go index 2d09e90f01..49f5818908 100644 --- a/pkg/cloudcommon/db/quotas/interface.go +++ b/pkg/cloudcommon/db/quotas/interface.go @@ -41,6 +41,7 @@ type IQuota interface { Update(quota IQuota) Add(quota IQuota) Sub(quota IQuota) + Allocable(quota IQuota) int ResetNegative() Exceed(request IQuota, quota IQuota) error // IsEmpty() bool @@ -68,6 +69,7 @@ type IQuotaManager interface { checkSetPendingQuota(ctx context.Context, userCred mcclient.TokenCredential, quota IQuota) error cancelPendingUsage(ctx context.Context, userCred mcclient.TokenCredential, localUsage IQuota, cancelUsage IQuota) error cancelUsage(ctx context.Context, userCred mcclient.TokenCredential, usage IQuota) error + getQuotaCount(ctx context.Context, request IQuota, pendingKey IQuotaKeys) (int, error) FetchIdNames(ctx context.Context, idMap map[string]map[string]string) (map[string]map[string]string, error) } diff --git a/pkg/cloudcommon/db/quotas/quotas.go b/pkg/cloudcommon/db/quotas/quotas.go index 4e46c1ee46..a56f0d2456 100644 --- a/pkg/cloudcommon/db/quotas/quotas.go +++ b/pkg/cloudcommon/db/quotas/quotas.go @@ -271,3 +271,49 @@ func (manager *SQuotaBaseManager) _checkSetPendingQuota(ctx context.Context, use } return nil } + +func (manager *SQuotaBaseManager) getQuotaCount(ctx context.Context, request IQuota, pendingKeys IQuotaKeys) (int, error) { + quotas, err := manager.GetParentQuotas(ctx, request.GetKeys()) + if err != nil { + return 0, errors.Wrap(err, "manager.getMatchedQuotas") + } + minCnt := -1 + for i := len(quotas) - 1; i >= 0; i -= 1 { + rel := relation(quotas[i].GetKeys(), pendingKeys) + if rel == QuotaKeysContain || rel == QuotaKeysEqual { + break + } + cnt, err := manager.__getQuotaCount(ctx, quotas[i], request) + if err != nil { + return 0, errors.Wrapf(err, "manager.__getQuotaCount for key %s", QuotaKeyString(quotas[i].GetKeys())) + } + if minCnt < 0 || minCnt > cnt { + minCnt = cnt + } + } + return minCnt, nil +} + +func (manager *SQuotaBaseManager) __getQuotaCount(ctx context.Context, quota IQuota, request IQuota) (int, error) { + keys := quota.GetKeys() + log.Debugf("__checkQuota for keys: %s", QuotaKeyString(keys)) + used := manager.newQuota() + err := manager.usageStore.GetQuota(ctx, keys, used) + if err != nil { + return 0, errors.Wrap(err, "manager.usageStore.GetQuotaByKeys") + } + log.Debugf("__checkQuota usage: %s", jsonutils.Marshal(used)) + pendings, err := manager.pendingStore.GetChildrenQuotas(ctx, keys) + if err != nil { + return 0, errors.Wrap(err, "manager.pendingStore.GetChildrenQuotas") + } + for i := range pendings { + if pendings[i].IsEmpty() { + continue + } + log.Debugf("__checkQuota pending %d: %s", i, jsonutils.Marshal(pendings[i])) + used.Add(pendings[i]) + } + quota.Sub(used) + return quota.Allocable(request), nil +} diff --git a/pkg/cloudcommon/db/quotas/register.go b/pkg/cloudcommon/db/quotas/register.go index ddf986e3c8..6e7c21dbbe 100644 --- a/pkg/cloudcommon/db/quotas/register.go +++ b/pkg/cloudcommon/db/quotas/register.go @@ -75,3 +75,8 @@ func cancelUsage(ctx context.Context, userCred mcclient.TokenCredential, usage I log.Errorf("cancelUsage %s fail: %s", jsonutils.Marshal(usage), err) } } + +func GetQuotaCount(ctx context.Context, request IQuota, pendingKeys IQuotaKeys) (int, error) { + manager := getQuotaManager(request) + return manager.getQuotaCount(ctx, request, pendingKeys) +} diff --git a/pkg/compute/models/cloudaccounts.go b/pkg/compute/models/cloudaccounts.go index 4d8676e4d1..e78183bce9 100644 --- a/pkg/compute/models/cloudaccounts.go +++ b/pkg/compute/models/cloudaccounts.go @@ -855,7 +855,7 @@ func (self *SCloudaccount) getProjectIds() []string { return ret } -func (self *SCloudaccount) getCloudEnv() string { +func (self *SCloudaccount) GetCloudEnv() string { if self.IsOnPremise { return api.CLOUD_ENV_ON_PREMISE } else if self.IsPublicCloud.IsTrue() { @@ -885,7 +885,7 @@ func (self *SCloudaccount) getMoreDetails(extra *jsonutils.JSONDict) *jsonutils. extra.Add(projects, "projects") extra.Set("sync_interval_seconds", jsonutils.NewInt(int64(self.getSyncIntervalSeconds()))) extra.Set("sync_status2", jsonutils.NewString(self.getSyncStatus2())) - extra.Set("cloud_env", jsonutils.NewString(self.getCloudEnv())) + extra.Set("cloud_env", jsonutils.NewString(self.GetCloudEnv())) if len(self.ProjectId) > 0 { if proj, _ := db.TenantCacheManager.FetchTenantById(context.Background(), self.ProjectId); proj != nil { extra.Add(jsonutils.NewString(proj.Name), "tenant") diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 741ce4b74d..96e66747b1 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -2091,7 +2091,7 @@ func (self *SGuest) PerformAttachnetwork(ctx context.Context, userCred mcclient. return nil, err } var inicCnt, enicCnt, ibw, ebw int - if isExitNetworkInfo(conf) { + if IsExitNetworkInfo(conf) { enicCnt = 1 ebw = conf.BwLimit } else { diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index fddaefddd8..752767b396 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -1467,7 +1467,7 @@ func getGuestResourceRequirements( eBw := 0 iBw := 0 for _, netConfig := range input.Networks { - if isExitNetworkInfo(netConfig) { + if IsExitNetworkInfo(netConfig) { eNicCnt += 1 eBw += netConfig.BwLimit } else { diff --git a/pkg/compute/models/managedresource.go b/pkg/compute/models/managedresource.go index e4055c0229..dd12344b13 100644 --- a/pkg/compute/models/managedresource.go +++ b/pkg/compute/models/managedresource.go @@ -550,7 +550,7 @@ func MakeCloudProviderInfoV2(region *SCloudregion, zone *SZone, provider *SCloud info.Provider = provider.Provider info.Brand = account.Brand - info.CloudEnv = account.getCloudEnv() + info.CloudEnv = account.GetCloudEnv() if region != nil { info.RegionExternalId = region.ExternalId @@ -602,7 +602,7 @@ func MakeCloudProviderInfo(region *SCloudregion, zone *SZone, provider *SCloudpr info.Account = account.GetName() info.AccountId = account.GetId() info.Brand = account.Brand - info.CloudEnv = account.getCloudEnv() + info.CloudEnv = account.GetCloudEnv() } if region != nil { diff --git a/pkg/compute/models/networks.go b/pkg/compute/models/networks.go index b600dc9943..679504761d 100644 --- a/pkg/compute/models/networks.go +++ b/pkg/compute/models/networks.go @@ -922,7 +922,7 @@ func isValidNetworkInfo(userCred mcclient.TokenCredential, netConfig *api.Networ return nil } -func isExitNetworkInfo(netConfig *api.NetworkConfig) bool { +func IsExitNetworkInfo(netConfig *api.NetworkConfig) bool { if len(netConfig.Network) > 0 { netObj, _ := NetworkManager.FetchById(netConfig.Network) net := netObj.(*SNetwork) diff --git a/pkg/compute/models/projectquota.go b/pkg/compute/models/projectquota.go index 5eabad0972..024470de76 100644 --- a/pkg/compute/models/projectquota.go +++ b/pkg/compute/models/projectquota.go @@ -141,6 +141,15 @@ func (self *SProjectQuota) Sub(quota quotas.IQuota) { self.Secgroup = nonNegative(self.Secgroup - squota.Secgroup) } +func (self *SProjectQuota) Allocable(request quotas.IQuota) int { + squota := request.(*SProjectQuota) + cnt := -1 + if self.Secgroup >= 0 && squota.Secgroup > 0 && (cnt < 0 || cnt > self.Secgroup/squota.Secgroup) { + cnt = self.Secgroup / squota.Secgroup + } + return cnt +} + func (self *SProjectQuota) Update(quota quotas.IQuota) { squota := quota.(*SProjectQuota) if squota.Secgroup > 0 { diff --git a/pkg/compute/models/quotas.go b/pkg/compute/models/quotas.go index f2e5c8f112..f0dc22e58e 100644 --- a/pkg/compute/models/quotas.go +++ b/pkg/compute/models/quotas.go @@ -258,6 +258,30 @@ func (self *SQuota) Sub(quota quotas.IQuota) { self.IsolatedDevice = nonNegative(self.IsolatedDevice - squota.IsolatedDevice) } +func (self *SQuota) Allocable(request quotas.IQuota) int { + squota := request.(*SQuota) + cnt := -1 + if self.Count >= 0 && squota.Count > 0 && (cnt < 0 || cnt > self.Count/squota.Count) { + cnt = self.Count / squota.Count + } + if self.Cpu >= 0 && squota.Cpu > 0 && (cnt < 0 || cnt > self.Cpu/squota.Cpu) { + cnt = self.Cpu / squota.Cpu + } + if self.Memory >= 0 && squota.Memory > 0 && (cnt < 0 || cnt > self.Memory/squota.Memory) { + cnt = self.Memory / squota.Memory + } + if self.Storage >= 0 && squota.Storage > 0 && (cnt < 0 || cnt > self.Storage/squota.Storage) { + cnt = self.Storage / squota.Storage + } + if self.Group >= 0 && squota.Group > 0 && (cnt < 0 || cnt > self.Group/squota.Group) { + cnt = self.Group / squota.Group + } + if self.IsolatedDevice >= 0 && squota.IsolatedDevice > 0 && (cnt < 0 || cnt > self.IsolatedDevice/squota.IsolatedDevice) { + cnt = self.IsolatedDevice / squota.IsolatedDevice + } + return cnt +} + func (self *SQuota) Update(quota quotas.IQuota) { squota := quota.(*SQuota) if squota.Count > 0 { @@ -446,7 +470,7 @@ func fetchCloudQuotaKeys(scope rbacutils.TRbacScope, ownerId mcclient.IIdentityP account := manager.GetCloudaccount() keys.Provider = account.Provider keys.Brand = account.Brand - keys.CloudEnv = account.getCloudEnv() + keys.CloudEnv = account.GetCloudEnv() keys.AccountId = account.Id keys.ManagerId = manager.Id } else { diff --git a/pkg/compute/models/regionquota.go b/pkg/compute/models/regionquota.go index 21ef5ea5b1..e05b2a66ae 100644 --- a/pkg/compute/models/regionquota.go +++ b/pkg/compute/models/regionquota.go @@ -312,6 +312,45 @@ func (self *SRegionQuota) Sub(quota quotas.IQuota) { self.Loadbalancer = nonNegative(self.Loadbalancer - squota.Loadbalancer) } +func (self *SRegionQuota) Allocable(request quotas.IQuota) int { + squota := request.(*SRegionQuota) + cnt := -1 + if self.Port >= 0 && squota.Port > 0 && (cnt < 0 || cnt > self.Port/squota.Port) { + cnt = self.Port / squota.Port + } + if self.Eip >= 0 && squota.Eip > 0 && (cnt < 0 || cnt > self.Eip/squota.Eip) { + cnt = self.Eip / squota.Eip + } + if self.Eport >= 0 && squota.Eport > 0 && (cnt < 0 || cnt > self.Eport/squota.Eport) { + cnt = self.Eport / squota.Eport + } + if self.Bw >= 0 && squota.Bw > 0 && (cnt < 0 || cnt > self.Bw/squota.Bw) { + cnt = self.Bw / squota.Bw + } + if self.Ebw >= 0 && squota.Ebw > 0 && (cnt < 0 || cnt > self.Ebw/squota.Ebw) { + cnt = self.Ebw / squota.Ebw + } + if self.Snapshot >= 0 && squota.Snapshot > 0 && (cnt < 0 || cnt > self.Snapshot/squota.Snapshot) { + cnt = self.Snapshot / squota.Snapshot + } + if self.Bucket >= 0 && squota.Bucket > 0 && (cnt < 0 || cnt > self.Bucket/squota.Bucket) { + cnt = self.Bucket / squota.Bucket + } + if self.ObjectGB >= 0 && squota.ObjectGB > 0 && (cnt < 0 || cnt > self.ObjectGB/squota.ObjectGB) { + cnt = self.ObjectGB / squota.ObjectGB + } + if self.ObjectCnt >= 0 && squota.ObjectCnt > 0 && (cnt < 0 || cnt > self.ObjectCnt/squota.ObjectCnt) { + cnt = self.ObjectCnt / squota.ObjectCnt + } + if self.Rds >= 0 && squota.Rds > 0 && (cnt < 0 || cnt > self.Rds/squota.Rds) { + cnt = self.Rds / squota.Rds + } + if self.Cache >= 0 && squota.Cache > 0 && (cnt < 0 || cnt > self.Cache/squota.Cache) { + cnt = self.Cache / squota.Cache + } + return cnt +} + func (self *SRegionQuota) Update(quota quotas.IQuota) { squota := quota.(*SRegionQuota) if squota.Port > 0 { diff --git a/pkg/compute/models/zonequota.go b/pkg/compute/models/zonequota.go index d3990cc18c..a1e1bcb3d4 100644 --- a/pkg/compute/models/zonequota.go +++ b/pkg/compute/models/zonequota.go @@ -174,6 +174,10 @@ func (self *SZoneQuota) Sub(quota quotas.IQuota) { // self.Loadbalancer = nonNegative(self.Loadbalancer - squota.Loadbalancer) } +func (self *SZoneQuota) Allocable(request quotas.IQuota) int { + return -1 +} + func (self *SZoneQuota) Update(quota quotas.IQuota) { // squota := quota.(*SZoneQuota) // if squota.Loadbalancer > 0 { diff --git a/pkg/compute/tasks/disk_batch_create_task.go b/pkg/compute/tasks/disk_batch_create_task.go index bab41f71f8..ec2f05bcd9 100644 --- a/pkg/compute/tasks/disk_batch_create_task.go +++ b/pkg/compute/tasks/disk_batch_create_task.go @@ -16,6 +16,7 @@ package tasks import ( "context" + "fmt" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -50,6 +51,7 @@ func (self *DiskBatchCreateTask) getNeedScheduleDisks(objs []db.IStandaloneModel func (self *DiskBatchCreateTask) clearPendingUsage(ctx context.Context, disk *models.SDisk) { ClearTaskPendingUsage(ctx, self) + ClearTaskPendingRegionUsage(ctx, self) } func (self *DiskBatchCreateTask) OnInit(ctx context.Context, objs []db.IStandaloneModel, body jsonutils.JSONObject) { @@ -82,6 +84,25 @@ func (self *DiskBatchCreateTask) GetSchedParams() (*schedapi.ScheduleInput, erro return ret, err } +func (self *DiskBatchCreateTask) GetDisks() ([]*api.DiskConfig, error) { + input, err := self.GetSchedParams() + if err != nil { + return nil, err + } + return input.Disks, nil +} + +func (self *DiskBatchCreateTask) GetFirstDisk() (*api.DiskConfig, error) { + disks, err := self.GetDisks() + if err != nil { + return nil, err + } + if len(disks) == 0 { + return nil, fmt.Errorf("Empty disks to schedule") + } + return disks[0], nil +} + func (self *DiskBatchCreateTask) OnScheduleFailCallback(ctx context.Context, obj IScheduleModel, reason string) { self.SSchedTask.OnScheduleFailCallback(ctx, obj, reason) disk := obj.(*models.SDisk) diff --git a/pkg/compute/tasks/guest_batch_create_task.go b/pkg/compute/tasks/guest_batch_create_task.go index f8fcdc2cf1..2733000149 100644 --- a/pkg/compute/tasks/guest_batch_create_task.go +++ b/pkg/compute/tasks/guest_batch_create_task.go @@ -16,12 +16,14 @@ package tasks import ( "context" + "fmt" "yunion.io/x/jsonutils" "yunion.io/x/log" api "yunion.io/x/onecloud/pkg/apis/compute" schedapi "yunion.io/x/onecloud/pkg/apis/scheduler" + "yunion.io/x/onecloud/pkg/cloudcommon/cmdline" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" @@ -39,6 +41,34 @@ func init() { taskman.RegisterTask(GuestBatchCreateTask{}) } +func (self *GuestBatchCreateTask) GetSchedParams() (*schedapi.ScheduleInput, error) { + params := self.GetParams() + input, err := cmdline.FetchScheduleInputByJSON(params) + if err != nil { + return nil, fmt.Errorf("Unmarsh to schedule input: %v", err) + } + return input, err +} + +func (self *GuestBatchCreateTask) GetDisks() ([]*api.DiskConfig, error) { + input, err := self.GetSchedParams() + if err != nil { + return nil, err + } + return input.Disks, nil +} + +func (self *GuestBatchCreateTask) GetFirstDisk() (*api.DiskConfig, error) { + disks, err := self.GetDisks() + if err != nil { + return nil, err + } + if len(disks) == 0 { + return nil, fmt.Errorf("Empty disks to schedule") + } + return disks[0], nil +} + func (self *GuestBatchCreateTask) GetCreateInput() (*api.ServerCreateInput, error) { input := new(api.ServerCreateInput) err := self.GetParams().Unmarshal(input) diff --git a/pkg/compute/tasks/pending_usage.go b/pkg/compute/tasks/pending_usage.go index 62be8a2bc9..2e2044e8c1 100644 --- a/pkg/compute/tasks/pending_usage.go +++ b/pkg/compute/tasks/pending_usage.go @@ -31,7 +31,13 @@ func ClearTaskPendingUsage(ctx context.Context, task taskman.ITask) error { err := task.GetPendingUsage(&pendingUsage, index) if err != nil { log.Errorf("GetPendingUsage fail %s", err) - return errors.Wrap(err, "task.GetPendingUsage") + // ignore error + // return errors.Wrap(err, "task.GetPendingUsage") + return nil + } + + if pendingUsage.IsEmpty() { + return nil } err = quotas.CancelPendingUsage(ctx, task.GetUserCred(), &pendingUsage, &pendingUsage) @@ -55,7 +61,13 @@ func ClearTaskPendingRegionUsage(ctx context.Context, task taskman.ITask) error err := task.GetPendingUsage(&pendingUsage, index) if err != nil { log.Errorf("GetPendingUsage fail %s", err) - return errors.Wrap(err, "task.GetPendingUsage") + // ignore error + // return errors.Wrap(err, "task.GetPendingUsage") + return nil + } + + if pendingUsage.IsEmpty() { + return nil } err = quotas.CancelPendingUsage(ctx, task.GetUserCred(), &pendingUsage, &pendingUsage) diff --git a/pkg/compute/tasks/schedule.go b/pkg/compute/tasks/schedule.go index 939ecef3a9..1bc7f4c784 100644 --- a/pkg/compute/tasks/schedule.go +++ b/pkg/compute/tasks/schedule.go @@ -19,11 +19,9 @@ import ( "fmt" "yunion.io/x/jsonutils" - "yunion.io/x/log" api "yunion.io/x/onecloud/pkg/apis/compute" schedapi "yunion.io/x/onecloud/pkg/apis/scheduler" - "yunion.io/x/onecloud/pkg/cloudcommon/cmdline" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" @@ -63,34 +61,6 @@ type SSchedTask struct { input *schedapi.ScheduleInput } -func (self *SSchedTask) GetSchedParams() (*schedapi.ScheduleInput, error) { - params := self.GetParams() - input, err := cmdline.FetchScheduleInputByJSON(params) - if err != nil { - return nil, fmt.Errorf("Unmarsh to schedule input: %v", err) - } - return input, err -} - -func (self *SSchedTask) GetDisks() ([]*api.DiskConfig, error) { - input, err := self.GetSchedParams() - if err != nil { - return nil, err - } - return input.Disks, nil -} - -func (self *SSchedTask) GetFirstDisk() (*api.DiskConfig, error) { - disks, err := self.GetDisks() - if err != nil { - return nil, err - } - if len(disks) == 0 { - return nil, fmt.Errorf("Empty disks to schedule") - } - return disks[0], nil -} - func (self *SSchedTask) OnStartSchedule(obj IScheduleModel) { db.OpsLog.LogEvent(obj, db.ACT_ALLOCATING, nil, self.GetUserCred()) obj.SetStatus(self.GetUserCred(), api.VM_SCHEDULE, "") @@ -145,6 +115,17 @@ func doScheduleObjects( } //schedInput = models.ApplySchedPolicies(schedInput) + // fetch pendingUsages + computeUsage := models.SQuota{} + task.GetPendingUsage(&computeUsage, 0) + regionUsage := models.SRegionQuota{} + task.GetPendingUsage(®ionUsage, 1) + + schedInput.PendingUsages = []jsonutils.JSONObject{ + jsonutils.Marshal(&computeUsage), + jsonutils.Marshal(®ionUsage), + } + params := jsonutils.Marshal(schedInput).(*jsonutils.JSONDict) task.SetStage("OnScheduleComplete", params) @@ -158,7 +139,9 @@ func doScheduleObjects( } func cancelPendingUsage(ctx context.Context, task IScheduleTask) { - pendingUsage := models.SQuota{} + ClearTaskPendingUsage(ctx, task.(taskman.ITask)) + ClearTaskPendingRegionUsage(ctx, task.(taskman.ITask)) + /*pendingUsage := models.SQuota{} err := task.GetPendingUsage(&pendingUsage, 0) if err != nil { log.Errorf("Taks GetPendingUsage fail %s", err) @@ -172,7 +155,7 @@ func cancelPendingUsage(ctx context.Context, task IScheduleTask) { } pendingRegionUsage := models.SRegionQuota{} - err = task.GetPendingUsage(&pendingRegionUsage, 0) + err = task.GetPendingUsage(&pendingRegionUsage, 1) if err != nil { log.Errorf("Taks GetRegionPendingUsage fail %s", err) return @@ -182,7 +165,7 @@ func cancelPendingUsage(ctx context.Context, task IScheduleTask) { if err != nil { log.Errorf("cancelpendingusage error %s", err) } - } + }*/ } func onSchedulerRequestFail( diff --git a/pkg/image/models/quotas.go b/pkg/image/models/quotas.go index 8dcec387b3..5c33d58683 100644 --- a/pkg/image/models/quotas.go +++ b/pkg/image/models/quotas.go @@ -155,6 +155,15 @@ func (self *SQuota) Sub(quota quotas.IQuota) { self.Image = quotas.NonNegative(self.Image - squota.Image) } +func (self *SQuota) Allocable(request quotas.IQuota) int { + squota := request.(*SQuota) + cnt := -1 + if self.Image >= 0 && squota.Image > 0 && (cnt < 0 || cnt > self.Image/squota.Image) { + cnt = self.Image / squota.Image + } + return cnt +} + func (self *SQuota) Update(quota quotas.IQuota) { squota := quota.(*SQuota) if squota.Image > 0 { diff --git a/pkg/scheduler/algorithm/predicates/quota_predicate.go b/pkg/scheduler/algorithm/predicates/quota_predicate.go new file mode 100644 index 0000000000..f62143a623 --- /dev/null +++ b/pkg/scheduler/algorithm/predicates/quota_predicate.go @@ -0,0 +1,131 @@ +// 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 predicates + +import ( + "context" + "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" + computemodels "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/scheduler/api" + "yunion.io/x/onecloud/pkg/scheduler/core" +) + +type SQuotaPredicate struct { + BasePredicate +} + +func (p *SQuotaPredicate) Name() string { + return "quota" +} + +func (p *SQuotaPredicate) Clone() core.FitPredicate { + return &SQuotaPredicate{} +} + +func (p *SQuotaPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) { + return true, nil +} + +func fetchGuestUsageFromSchedInfo(s *api.SchedInfo) (computemodels.SQuota, computemodels.SRegionQuota) { + vcpuCount := s.Ncpu + if vcpuCount == 0 { + vcpuCount = 1 + } + + vmemSize := s.Memory + + diskSize := 0 + + for _, diskConfig := range s.Disks { + diskSize += diskConfig.SizeMb + } + + devCount := len(s.IsolatedDevices) + + eNicCnt := 0 + iNicCnt := 0 + eBw := 0 + iBw := 0 + + for _, netConfig := range s.Networks { + if computemodels.IsExitNetworkInfo(netConfig) { + eNicCnt += 1 + eBw += netConfig.BwLimit + } else { + iNicCnt += 1 + iBw += netConfig.BwLimit + } + } + + req := computemodels.SQuota{ + Count: 1, + Cpu: int(vcpuCount), + Memory: int(vmemSize), + Storage: diskSize, + IsolatedDevice: devCount, + } + regionReq := computemodels.SRegionQuota{ + Port: iNicCnt, + Eport: eNicCnt, + Bw: iBw, + Ebw: eBw, + } + return req, regionReq +} + +func (p *SQuotaPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) { + h := NewPredicateHelper(p, u, c) + + d := u.SchedData() + + computePending := computemodels.SQuota{} + regionPending := computemodels.SRegionQuota{} + if len(d.PendingUsages) > 0 { + d.PendingUsages[0].Unmarshal(&computePending) + } + if len(d.PendingUsages) > 1 { + d.PendingUsages[1].Unmarshal(®ionPending) + } + + computeKeys := c.Getter().GetQuotaKeys(d) + + computeQuota, regionQuota := fetchGuestUsageFromSchedInfo(d) + + computeQuota.SetKeys(computeKeys) + regionQuota.SetKeys(computeKeys.SRegionalCloudResourceKeys) + + ctx := context.Background() + minCnt := -1 + if !computePending.IsEmpty() { + computeCnt, _ := quotas.GetQuotaCount(ctx, &computeQuota, computePending.GetKeys()) + if minCnt < 0 || minCnt > computeCnt { + minCnt = computeCnt + } + } + if !regionPending.IsEmpty() { + regionCnt, _ := quotas.GetQuotaCount(ctx, ®ionQuota, regionPending.GetKeys()) + if minCnt < 0 || minCnt > regionCnt { + minCnt = regionCnt + } + } + + if minCnt == 0 { + h.Exclude("quota limit") + } else if minCnt > 0 { + h.SetCapacity(int64(minCnt)) + } + + return h.GetResult() +} diff --git a/pkg/scheduler/algorithmprovider/defaults.go b/pkg/scheduler/algorithmprovider/defaults.go index 64450f1595..3ec00a7eb0 100644 --- a/pkg/scheduler/algorithmprovider/defaults.go +++ b/pkg/scheduler/algorithmprovider/defaults.go @@ -47,6 +47,7 @@ func defaultPredicates() sets.String { factory.RegisterFitPredicate("o-GuestNetschedtagFilter", &predicates.NetworkSchedtagPredicate{}), factory.RegisterFitPredicate("p-GuestForcedDispersionFilter", &predicates.SForcedGroupPredicate{}), factory.RegisterFitPredicate("p-GuestUnForcedDispersionFilter", &predicates.SUnForcedGroupPredicate{}), + factory.RegisterFitPredicate("z-QuotaFilter", &predicates.SQuotaPredicate{}), ) } diff --git a/pkg/scheduler/cache/candidate/base.go b/pkg/scheduler/cache/candidate/base.go index 92f84707e7..57fbc47163 100644 --- a/pkg/scheduler/cache/candidate/base.go +++ b/pkg/scheduler/cache/candidate/base.go @@ -35,6 +35,7 @@ type BaseHostDesc struct { Region *computemodels.SCloudregion `json:"region"` Zone *computemodels.SZone `json:"zone"` Cloudprovider *computemodels.SCloudprovider `json:"cloudprovider"` + Cloudaccount *computemodels.SCloudaccount `json:"cloudaccount"` Networks []*api.CandidateNetwork `json:"networks"` NetInterfaces map[string][]computemodels.SNetInterface `json:"net_interfaces"` Storages []*api.CandidateStorage `json:"storages"` @@ -197,6 +198,10 @@ func (b baseHostGetter) GetIpmiInfo() types.SIPMIInfo { return b.h.IpmiInfo } +func (b baseHostGetter) GetQuotaKeys(s *api.SchedInfo) computemodels.SComputeResourceKeys { + return b.h.getQuotaKeys(s) +} + func reviseResourceType(resType string) string { if resType == "" { return computeapi.HostResourceTypeDefault @@ -292,6 +297,7 @@ func (b BaseHostDesc) GetResourceType() string { func (b *BaseHostDesc) fillCloudProvider(host *computemodels.SHost) error { b.Cloudprovider = host.GetCloudprovider() + b.Cloudaccount = b.Cloudprovider.GetCloudaccount() return nil } @@ -417,6 +423,29 @@ func (h *BaseHostDesc) GetHostType() string { return h.HostType } +func (h *BaseHostDesc) getQuotaKeys(s *api.SchedInfo) computemodels.SComputeResourceKeys { + computeKeys := computemodels.SComputeResourceKeys{} + computeKeys.DomainId = s.Domain + computeKeys.ProjectId = s.Project + if h.Cloudprovider != nil { + computeKeys.Provider = h.Cloudaccount.Provider + computeKeys.Brand = h.Cloudaccount.Brand + computeKeys.CloudEnv = h.Cloudaccount.GetCloudEnv() + computeKeys.AccountId = h.Cloudaccount.Id + computeKeys.ManagerId = h.Cloudprovider.Id + } else { + computeKeys.Provider = computeapi.CLOUD_PROVIDER_ONECLOUD + computeKeys.Brand = computeapi.ONECLOUD_BRAND_ONECLOUD + computeKeys.CloudEnv = computeapi.CLOUD_ENV_ON_PREMISE + computeKeys.AccountId = "" + computeKeys.ManagerId = "" + } + computeKeys.RegionId = h.Region.Id + computeKeys.ZoneId = h.Zone.Id + computeKeys.Hypervisor = computeapi.HOSTTYPE_HYPERVISOR[h.HostType] + return computeKeys +} + func HostsResidentTenantStats(hostIDs []string) (map[string]map[string]interface{}, error) { residentTenantStats, err := FetchHostsResidentTenants(hostIDs) if err != nil { diff --git a/pkg/scheduler/core/types.go b/pkg/scheduler/core/types.go index 6842ce95e5..09a6c7502e 100644 --- a/pkg/scheduler/core/types.go +++ b/pkg/scheduler/core/types.go @@ -92,6 +92,8 @@ type CandidatePropertyGetter interface { GetFreeGroupCount(groupId string) (int, error) GetIpmiInfo() types.SIPMIInfo + + GetQuotaKeys(s *api.SchedInfo) computemodels.SComputeResourceKeys } // Candidater replace host Candidate resource info