From bc7f62a2add490be9a42385088d4829803ec70af Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 3 Jan 2020 00:40:00 +0800 Subject: [PATCH 1/6] 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 From 82f4885b3bf713a2502f1251e4a6ac28bd3e1be9 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 3 Jan 2020 00:54:31 +0800 Subject: [PATCH 2/6] move lb from zonequota to regionquota --- pkg/compute/models/regionquota.go | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pkg/compute/models/regionquota.go b/pkg/compute/models/regionquota.go index e05b2a66ae..6c7e47febc 100644 --- a/pkg/compute/models/regionquota.go +++ b/pkg/compute/models/regionquota.go @@ -348,6 +348,9 @@ func (self *SRegionQuota) Allocable(request quotas.IQuota) int { if self.Cache >= 0 && squota.Cache > 0 && (cnt < 0 || cnt > self.Cache/squota.Cache) { cnt = self.Cache / squota.Cache } + if self.Loadbalancer >= 0 && squota.Loadbalancer > 0 && (cnt < 0 || cnt > self.Loadbalancer/squota.Loadbalancer) { + cnt = self.Loadbalancer / squota.Loadbalancer + } return cnt } From 4a8a79a786cfdcec876e03520bfbc2fe43a1c247 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 3 Jan 2020 13:52:28 +0800 Subject: [PATCH 3/6] bug fixes --- cmd/climc/shell/quotas.go | 2 +- pkg/cloudcommon/db/db_dispatcher.go | 69 ++++++++++----- pkg/cloudcommon/db/interface.go | 9 +- pkg/cloudcommon/db/modelbase.go | 15 +++- pkg/cloudcommon/db/quotas/context.go | 87 +++++++++++++++++++ pkg/cloudcommon/db/quotas/quotas.go | 13 ++- pkg/cloudcommon/db/quotas/register.go | 2 + pkg/cloudcommon/db/usages.go | 34 ++++++++ pkg/compute/models/guestnetworks.go | 1 + pkg/compute/models/guests.go | 48 ++++------ pkg/compute/service/service.go | 3 + .../algorithm/predicates/quota_predicate.go | 9 +- pkg/scheduler/cache/candidate/base.go | 4 +- 13 files changed, 218 insertions(+), 78 deletions(-) create mode 100644 pkg/cloudcommon/db/quotas/context.go create mode 100644 pkg/cloudcommon/db/usages.go diff --git a/cmd/climc/shell/quotas.go b/cmd/climc/shell/quotas.go index 39772e1d3e..c4f342cb29 100644 --- a/cmd/climc/shell/quotas.go +++ b/cmd/climc/shell/quotas.go @@ -24,7 +24,7 @@ import ( ) type ComputeQuotaKeys struct { - RegionQuotaKeys + ZoneQuotaKeys Hypervisor string `help:"hypervisor" choices:"kvm|baremetal"` } diff --git a/pkg/cloudcommon/db/db_dispatcher.go b/pkg/cloudcommon/db/db_dispatcher.go index b30eefd801..c6fe724091 100644 --- a/pkg/cloudcommon/db/db_dispatcher.go +++ b/pkg/cloudcommon/db/db_dispatcher.go @@ -24,6 +24,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/gotypes" "yunion.io/x/pkg/util/filterclause" "yunion.io/x/pkg/utils" @@ -43,10 +44,6 @@ import ( "yunion.io/x/onecloud/pkg/util/stringutils2" ) -var ( - CancelUsages func(ctx context.Context, userCred mcclient.TokenCredential, usages []IUsage) -) - type DBModelDispatcher struct { modelManager IModelManager } @@ -1104,9 +1101,15 @@ func (dispatcher *DBModelDispatcher) Create(ctx context.Context, query jsonutils return nil, httperrors.NewForbiddenError("Not allow to create item") } + ctx = InitPendingUsagesInContext(ctx) + model, err := DoCreate(dispatcher.modelManager, ctx, userCred, query, data, ownerId) if err != nil { - log.Errorf("fail to doCreateItem %s", err) + // log.Errorf("fail to doCreateItem %s", err) + failErr := manager.OnCreateFailed(ctx, userCred, ownerId, query, data) + if failErr != nil { + log.Errorf("manager.OnCreateFailed %s", failErr) + } return nil, httperrors.NewGeneralError(err) } @@ -1179,41 +1182,61 @@ func (dispatcher *DBModelDispatcher) BatchCreate(ctx context.Context, query json } var ( - multiData []jsonutils.JSONObject - onBatchCreateFail func() - validateError error + multiData []jsonutils.JSONObject + // onBatchCreateFail func() + // validateError error ) + ctx = InitPendingUsagesInContext(ctx) + createResults, err := func() ([]sCreateResult, error) { lockman.LockClass(ctx, manager, GetLockClassKey(manager, ownerId)) defer lockman.ReleaseClass(ctx, manager, GetLockClassKey(manager, ownerId)) - multiData, err = expandMultiCreateParams(data, count) + // invoke only Once + err = manager.BatchPreValidate(ctx, userCred, ownerId, query, data.(*jsonutils.JSONDict), count) if err != nil { - return nil, err + return nil, errors.Wrap(err, "manager.BatchPreValidate") } + multiData, err = expandMultiCreateParams(data, count) + if err != nil { + return nil, errors.Wrap(err, "expandMultiCreateParams") + } + + // one fail, then all fail ret := make([]sCreateResult, len(multiData)) - for i, cdata := range multiData { - if i == 0 { - onBatchCreateFail, validateError = manager.BatchPreValidate( - ctx, userCred, ownerId, query, cdata.(*jsonutils.JSONDict), len(multiData)) - if validateError != nil { - return nil, validateError + for i := range multiData { + var model IModel + model, err = batchCreateDoCreateItem(manager, ctx, userCred, ownerId, query, multiData[i], i) + if err == nil { + ret[i] = sCreateResult{model: model, err: nil} + } else { + break + } + } + if err != nil { + for i := range ret { + if ret[i].model != nil { + DeleteModel(ctx, userCred, ret[i].model) + ret[i].model = nil } - } - model, err := batchCreateDoCreateItem(manager, ctx, userCred, ownerId, query, cdata, i) - if err != nil && onBatchCreateFail != nil { - onBatchCreateFail() + ret[i].err = err } - ret[i] = sCreateResult{model: model, err: err} + return nil, errors.Wrap(err, "batchCreateDoCreateItem") + } else { + return ret, nil } - return ret, nil }() if err != nil { - return nil, err + failErr := manager.OnCreateFailed(ctx, userCred, ownerId, query, data) + if failErr != nil { + log.Errorf("manager.OnCreateFailed %s", failErr) + } + return nil, errors.Wrap(err, "createResults") } + results := make([]modulebase.SubmitResult, count) models := make([]IModel, 0) for i, res := range createResults { diff --git a/pkg/cloudcommon/db/interface.go b/pkg/cloudcommon/db/interface.go index 0a817e93ab..6283bf8c62 100644 --- a/pkg/cloudcommon/db/interface.go +++ b/pkg/cloudcommon/db/interface.go @@ -30,11 +30,6 @@ import ( "yunion.io/x/onecloud/pkg/util/stringutils2" ) -type IUsage interface { - FetchUsage(ctx context.Context) error - IsEmpty() bool -} - type IModelManager interface { lockman.ILockedClass object.IObject @@ -91,7 +86,9 @@ type IModelManager interface { // ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) OnCreateComplete(ctx context.Context, items []IModel, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) BatchPreValidate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, - query jsonutils.JSONObject, data *jsonutils.JSONDict, count int) (func(), error) + query jsonutils.JSONObject, data *jsonutils.JSONDict, count int) error + + OnCreateFailed(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) error // allow perform action AllowPerformAction(ctx context.Context, userCred mcclient.TokenCredential, action string, query jsonutils.JSONObject, data jsonutils.JSONObject) bool diff --git a/pkg/cloudcommon/db/modelbase.go b/pkg/cloudcommon/db/modelbase.go index 3d7dedc384..9b5c7bbe39 100644 --- a/pkg/cloudcommon/db/modelbase.go +++ b/pkg/cloudcommon/db/modelbase.go @@ -22,6 +22,7 @@ import ( "time" "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" "yunion.io/x/sqlchemy" "yunion.io/x/onecloud/pkg/apis" @@ -381,14 +382,24 @@ func (manager *SModelBaseManager) GetPropertyDistinctField(ctx context.Context, func (manager *SModelBaseManager) BatchPreValidate( ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data *jsonutils.JSONDict, count int, -) (func(), error) { - return nil, nil +) error { + return nil } func (manager *SModelBaseManager) BatchCreateValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { return nil, nil } +func (manager *SModelBaseManager) OnCreateFailed(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) error { + if CancelPendingUsagesInContext != nil { + err := CancelPendingUsagesInContext(ctx, userCred) + if err != nil { + return errors.Wrap(err, "CancelPendingUsagesInContext") + } + } + return nil +} + func (model *SModelBase) GetId() string { return "" } diff --git a/pkg/cloudcommon/db/quotas/context.go b/pkg/cloudcommon/db/quotas/context.go new file mode 100644 index 0000000000..53e67479b8 --- /dev/null +++ b/pkg/cloudcommon/db/quotas/context.go @@ -0,0 +1,87 @@ +// 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 quotas + +import ( + "container/list" + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + + "yunion.io/x/onecloud/pkg/appctx" + "yunion.io/x/onecloud/pkg/mcclient" +) + +const ( + APP_CONTEXT_KEY_PENDINGUSAGES = appctx.AppContextKey("pendingusages") +) + +func initPendingUsagesInContext(ctx context.Context) context.Context { + return context.WithValue(ctx, APP_CONTEXT_KEY_PENDINGUSAGES, list.New()) +} + +func appContextPendingUsages(ctx context.Context) []IQuota { + val := ctx.Value(APP_CONTEXT_KEY_PENDINGUSAGES) + if val != nil { + quotaList := val.(*list.List) + ret := make([]IQuota, 0) + for e := quotaList.Front(); e != nil; e = e.Next() { + ret = append(ret, e.Value.(IQuota)) + } + return ret + } else { + return nil + } +} + +func clearPendingUsagesInContext(ctx context.Context) { + val := ctx.Value(APP_CONTEXT_KEY_PENDINGUSAGES) + if val != nil { + quotaList := val.(*list.List) + for quotaList.Len() > 0 { + quotaList.Remove(quotaList.Front()) + } + } +} + +func SavePendingUsagesInContext(ctx context.Context, quotas ...IQuota) { + val := ctx.Value(APP_CONTEXT_KEY_PENDINGUSAGES) + if val != nil { + quotaList := val.(*list.List) + for i := range quotas { + quotaList.PushBack(quotas[i]) + } + } +} + +func cancelPendingUsagesInContext(ctx context.Context, userCred mcclient.TokenCredential) error { + quotas := appContextPendingUsages(ctx) + if quotas == nil { + return nil + } + errs := make([]error, 0) + for i := range quotas { + err := CancelPendingUsage(ctx, userCred, quotas[i], quotas[i]) + if err != nil { + errs = append(errs, errors.Wrapf(err, "CancelPendingUsage %s", jsonutils.Marshal(quotas[i]))) + } + } + if len(errs) > 0 { + return errors.NewAggregate(errs) + } + clearPendingUsagesInContext(ctx) + return nil +} diff --git a/pkg/cloudcommon/db/quotas/quotas.go b/pkg/cloudcommon/db/quotas/quotas.go index a56f0d2456..dc96097ef4 100644 --- a/pkg/cloudcommon/db/quotas/quotas.go +++ b/pkg/cloudcommon/db/quotas/quotas.go @@ -216,13 +216,11 @@ func (manager *SQuotaBaseManager) checkQuota(ctx context.Context, request IQuota func (manager *SQuotaBaseManager) __checkQuota(ctx context.Context, quota IQuota, request IQuota) 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 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 errors.Wrap(err, "manager.pendingStore.GetChildrenQuotas") @@ -231,7 +229,6 @@ func (manager *SQuotaBaseManager) __checkQuota(ctx context.Context, quota IQuota if pendings[i].IsEmpty() { continue } - log.Debugf("__checkQuota pending %d: %s", i, jsonutils.Marshal(pendings[i])) used.Add(pendings[i]) } return used.Exceed(request, quota) @@ -296,13 +293,11 @@ func (manager *SQuotaBaseManager) getQuotaCount(ctx context.Context, request IQu 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") @@ -311,9 +306,13 @@ func (manager *SQuotaBaseManager) __getQuotaCount(ctx context.Context, quota IQu if pendings[i].IsEmpty() { continue } - log.Debugf("__checkQuota pending %d: %s", i, jsonutils.Marshal(pendings[i])) used.Add(pendings[i]) } + err = used.Exceed(request, quota) + if err != nil { + return 0, nil + } quota.Sub(used) - return quota.Allocable(request), nil + cnt := quota.Allocable(request) + return cnt, nil } diff --git a/pkg/cloudcommon/db/quotas/register.go b/pkg/cloudcommon/db/quotas/register.go index 6e7c21dbbe..5ac0944c3c 100644 --- a/pkg/cloudcommon/db/quotas/register.go +++ b/pkg/cloudcommon/db/quotas/register.go @@ -33,6 +33,8 @@ func init() { quotaManagerTable = make(map[reflect.Type]IQuotaManager) db.CancelUsages = CancelUsages + db.CancelPendingUsagesInContext = cancelPendingUsagesInContext + db.InitPendingUsagesInContext = initPendingUsagesInContext } func Register(manager IQuotaManager) { diff --git a/pkg/cloudcommon/db/usages.go b/pkg/cloudcommon/db/usages.go new file mode 100644 index 0000000000..65f31b061f --- /dev/null +++ b/pkg/cloudcommon/db/usages.go @@ -0,0 +1,34 @@ +// 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 db + +import ( + "context" + + "yunion.io/x/onecloud/pkg/mcclient" +) + +type IUsage interface { + FetchUsage(ctx context.Context) error + IsEmpty() bool +} + +var ( + CancelUsages func(ctx context.Context, userCred mcclient.TokenCredential, usages []IUsage) + + CancelPendingUsagesInContext func(ctx context.Context, userCred mcclient.TokenCredential) error + + InitPendingUsagesInContext func(ctx context.Context) context.Context +) diff --git a/pkg/compute/models/guestnetworks.go b/pkg/compute/models/guestnetworks.go index bd67cad261..0520c3acee 100644 --- a/pkg/compute/models/guestnetworks.go +++ b/pkg/compute/models/guestnetworks.go @@ -540,6 +540,7 @@ func totalGuestNicCount( guests := GuestManager.Query().SubQuery() hosts := HostManager.Query().SubQuery() guestnics := GuestnetworkManager.Query().SubQuery() + q := guestnics.Query() q = q.Join(guests, sqlchemy.Equals(guests.Field("id"), guestnics.Field("guest_id"))) q = q.Join(hosts, sqlchemy.Equals(guests.Field("host_id"), hosts.Field("id"))) diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 752767b396..ad6a844ef8 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -904,40 +904,18 @@ func serverCreateInput2ComputeQuotaKeys(input api.ServerCreateInput, ownerId mcc func (manager *SGuestManager) BatchPreValidate( ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data *jsonutils.JSONDict, count int, -) (func(), error) { +) error { input, err := manager.validateCreateData(ctx, userCred, ownerId, query, data) if err != nil { - return nil, err + return errors.Wrap(err, "manager.validateCreateData") } - if input.IsSystem == nil || *input.IsSystem == false { - reqQuota, reqRegionQuota, err := manager.checkCreateQuota(ctx, userCred, ownerId, *input, input.Backup, count) + if input.IsSystem == nil || !(*input.IsSystem) { + err := manager.checkCreateQuota(ctx, userCred, ownerId, *input, input.Backup, count) if err != nil { - return nil, err + return errors.Wrap(err, "manager.checkCreateQuota") } - quota := &SQuota{ - Count: reqQuota.Count / count, - Cpu: reqQuota.Cpu / count, - Memory: reqQuota.Memory / count, - Storage: reqQuota.Storage / count, - IsolatedDevice: reqQuota.IsolatedDevice / count, - } - regionQuota := &SRegionQuota{ - Port: reqRegionQuota.Port / count, - Eport: reqRegionQuota.Eport / count, - Bw: reqRegionQuota.Bw / count, - Ebw: reqRegionQuota.Ebw / count, - Eip: reqRegionQuota.Eip / count, - } - keys := serverCreateInput2ComputeQuotaKeys(*input, ownerId) - regionKeys := keys.SRegionalCloudResourceKeys - quota.SetKeys(keys) - regionQuota.SetKeys(regionKeys) - return func() { - quotas.CancelPendingUsage(ctx, userCred, quota, quota) - quotas.CancelPendingUsage(ctx, userCred, regionQuota, regionQuota) - }, nil } - return nil, nil + return nil } func parseInstanceSnapshot(input *api.ServerCreateInput) (*api.ServerCreateInput, error) { @@ -1335,7 +1313,7 @@ func (manager *SGuestManager) ValidateCreateData(ctx context.Context, userCred m return nil, err } if input.IsSystem == nil || !(*input.IsSystem) { - _, _, err = manager.checkCreateQuota(ctx, userCred, ownerId, *input, input.Backup, 1) + err = manager.checkCreateQuota(ctx, userCred, ownerId, *input, input.Backup, 1) if err != nil { return nil, err } @@ -1402,17 +1380,21 @@ func (manager *SGuestManager) checkCreateQuota( input api.ServerCreateInput, hasBackup bool, count int, -) (*SQuota, *SRegionQuota, error) { +) error { req, regionReq := getGuestResourceRequirements(ctx, userCred, input, ownerId, count, hasBackup) + log.Debugf("computeQuota: %s", jsonutils.Marshal(req)) + log.Debugf("regionQuota: %s", jsonutils.Marshal(regionReq)) + err := quotas.CheckSetPendingQuota(ctx, userCred, &req) if err != nil { - return nil, nil, err + return errors.Wrap(err, "quotas.CheckSetPendingQuota") } err = quotas.CheckSetPendingQuota(ctx, userCred, ®ionReq) if err != nil { - return nil, nil, err + return errors.Wrap(err, "quotas.CheckSetPendingQuota") } - return &req, ®ionReq, nil + quotas.SavePendingUsagesInContext(ctx, &req, ®ionReq) + return nil } func (self *SGuest) checkUpdateQuota(ctx context.Context, userCred mcclient.TokenCredential, vcpuCount int, vmemSize int) (quotas.IQuota, error) { diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index ca9b2e84c2..0241dc381d 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -84,6 +84,9 @@ func StartService() { cron.AddJobAtIntervals("StartHostPingDetectionTask", time.Duration(opts.HostOfflineDetectionInterval)*time.Second, models.HostManager.PingDetectionTask) cron.AddJobAtIntervalsWithStartRun("CalculateQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.QuotaManager.CalculateQuotaUsages, true) + cron.AddJobAtIntervalsWithStartRun("CalculateRegionQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.RegionQuotaManager.CalculateQuotaUsages, true) + cron.AddJobAtIntervalsWithStartRun("CalculateZoneQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.ZoneQuotaManager.CalculateQuotaUsages, true) + cron.AddJobAtIntervalsWithStartRun("CalculateProjectQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.ProjectQuotaManager.CalculateQuotaUsages, true) cron.AddJobAtIntervalsWithStartRun("AutoSyncCloudaccountTask", time.Duration(opts.CloudAutoSyncIntervalSeconds)*time.Second, models.CloudaccountManager.AutoSyncCloudaccountTask, true) diff --git a/pkg/scheduler/algorithm/predicates/quota_predicate.go b/pkg/scheduler/algorithm/predicates/quota_predicate.go index f62143a623..2a3a62b1f0 100644 --- a/pkg/scheduler/algorithm/predicates/quota_predicate.go +++ b/pkg/scheduler/algorithm/predicates/quota_predicate.go @@ -16,6 +16,7 @@ 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" @@ -110,20 +111,18 @@ func (p *SQuotaPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core minCnt := -1 if !computePending.IsEmpty() { computeCnt, _ := quotas.GetQuotaCount(ctx, &computeQuota, computePending.GetKeys()) - if minCnt < 0 || minCnt > computeCnt { + if computeCnt >= 0 && (minCnt < 0 || minCnt > computeCnt) { minCnt = computeCnt } } if !regionPending.IsEmpty() { regionCnt, _ := quotas.GetQuotaCount(ctx, ®ionQuota, regionPending.GetKeys()) - if minCnt < 0 || minCnt > regionCnt { + if regionCnt >= 0 && (minCnt < 0 || minCnt > regionCnt) { minCnt = regionCnt } } - if minCnt == 0 { - h.Exclude("quota limit") - } else if minCnt > 0 { + if minCnt >= 0 { h.SetCapacity(int64(minCnt)) } diff --git a/pkg/scheduler/cache/candidate/base.go b/pkg/scheduler/cache/candidate/base.go index 57fbc47163..cb0f0587be 100644 --- a/pkg/scheduler/cache/candidate/base.go +++ b/pkg/scheduler/cache/candidate/base.go @@ -297,7 +297,9 @@ func (b BaseHostDesc) GetResourceType() string { func (b *BaseHostDesc) fillCloudProvider(host *computemodels.SHost) error { b.Cloudprovider = host.GetCloudprovider() - b.Cloudaccount = b.Cloudprovider.GetCloudaccount() + if b.Cloudprovider != nil { + b.Cloudaccount = b.Cloudprovider.GetCloudaccount() + } return nil } From f62cea0e1d533c33e75e968927a5e8da69ccd8a4 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 3 Jan 2020 13:57:24 +0800 Subject: [PATCH 4/6] remove obsolete codes --- pkg/compute/tasks/schedule.go | 25 ------------------------- 1 file changed, 25 deletions(-) diff --git a/pkg/compute/tasks/schedule.go b/pkg/compute/tasks/schedule.go index 1bc7f4c784..93ece147b5 100644 --- a/pkg/compute/tasks/schedule.go +++ b/pkg/compute/tasks/schedule.go @@ -141,31 +141,6 @@ func doScheduleObjects( func cancelPendingUsage(ctx context.Context, task IScheduleTask) { 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) - return - } - if !pendingUsage.IsEmpty() { - err = quotas.CancelPendingUsage(ctx, task.GetUserCred(), &pendingUsage, &pendingUsage) - if err != nil { - log.Errorf("cancelpendingusage error %s", err) - } - } - - pendingRegionUsage := models.SRegionQuota{} - err = task.GetPendingUsage(&pendingRegionUsage, 1) - if err != nil { - log.Errorf("Taks GetRegionPendingUsage fail %s", err) - return - } - if !pendingRegionUsage.IsEmpty() { - err = quotas.CancelPendingUsage(ctx, task.GetUserCred(), &pendingRegionUsage, &pendingRegionUsage) - if err != nil { - log.Errorf("cancelpendingusage error %s", err) - } - }*/ } func onSchedulerRequestFail( From 940f6c84dee569505133610fc12ce4437b5428b1 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 3 Jan 2020 15:18:19 +0800 Subject: [PATCH 5/6] automatically save pendingUsage in context --- pkg/cloudcommon/db/db_dispatcher.go | 12 +++--------- pkg/cloudcommon/db/quotas/context.go | 2 +- pkg/cloudcommon/db/quotas/register.go | 8 +++++++- pkg/compute/models/guests.go | 1 - 4 files changed, 11 insertions(+), 12 deletions(-) diff --git a/pkg/cloudcommon/db/db_dispatcher.go b/pkg/cloudcommon/db/db_dispatcher.go index c6fe724091..c7b5660912 100644 --- a/pkg/cloudcommon/db/db_dispatcher.go +++ b/pkg/cloudcommon/db/db_dispatcher.go @@ -39,7 +39,6 @@ import ( "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/auth" "yunion.io/x/onecloud/pkg/mcclient/modulebase" - "yunion.io/x/onecloud/pkg/util/httputils" "yunion.io/x/onecloud/pkg/util/rbacutils" "yunion.io/x/onecloud/pkg/util/stringutils2" ) @@ -1242,14 +1241,9 @@ func (dispatcher *DBModelDispatcher) BatchCreate(ctx context.Context, query json for i, res := range createResults { result := modulebase.SubmitResult{} if res.err != nil { - jsonErr, ok := res.err.(*httputils.JSONClientError) - if ok { - result.Status = jsonErr.Code - result.Data = jsonutils.Marshal(jsonErr) - } else { - result.Status = 500 - result.Data = jsonutils.NewString(res.err.Error()) - } + jsonErr := httperrors.NewGeneralError(res.err) + result.Status = jsonErr.Code + result.Data = jsonutils.Marshal(jsonErr) } else { lockman.LockObject(ctx, res.model) defer lockman.ReleaseObject(ctx, res.model) diff --git a/pkg/cloudcommon/db/quotas/context.go b/pkg/cloudcommon/db/quotas/context.go index 53e67479b8..2da99b30fc 100644 --- a/pkg/cloudcommon/db/quotas/context.go +++ b/pkg/cloudcommon/db/quotas/context.go @@ -57,7 +57,7 @@ func clearPendingUsagesInContext(ctx context.Context) { } } -func SavePendingUsagesInContext(ctx context.Context, quotas ...IQuota) { +func savePendingUsagesInContext(ctx context.Context, quotas ...IQuota) { val := ctx.Value(APP_CONTEXT_KEY_PENDINGUSAGES) if val != nil { quotaList := val.(*list.List) diff --git a/pkg/cloudcommon/db/quotas/register.go b/pkg/cloudcommon/db/quotas/register.go index 5ac0944c3c..80ef522e81 100644 --- a/pkg/cloudcommon/db/quotas/register.go +++ b/pkg/cloudcommon/db/quotas/register.go @@ -20,6 +20,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/mcclient" @@ -61,7 +62,12 @@ func CancelPendingUsage(ctx context.Context, userCred mcclient.TokenCredential, func CheckSetPendingQuota(ctx context.Context, userCred mcclient.TokenCredential, quota IQuota) error { manager := getQuotaManager(quota) - return manager.checkSetPendingQuota(ctx, userCred, quota) + err := manager.checkSetPendingQuota(ctx, userCred, quota) + if err != nil { + return errors.Wrap(err, "manager.checkSetPendingQuota") + } + savePendingUsagesInContext(ctx, quota) + return nil } func CancelUsages(ctx context.Context, userCred mcclient.TokenCredential, usages []db.IUsage) { diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index ad6a844ef8..b707b383fd 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -1393,7 +1393,6 @@ func (manager *SGuestManager) checkCreateQuota( if err != nil { return errors.Wrap(err, "quotas.CheckSetPendingQuota") } - quotas.SavePendingUsagesInContext(ctx, &req, ®ionReq) return nil } From 99bb6af52c243ef4ece9b727b0812ec5d835546e Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Fri, 3 Jan 2020 16:00:11 +0800 Subject: [PATCH 6/6] fix: OutOfQuota error code 500 --- pkg/mcclient/modulebase/resource.go | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/pkg/mcclient/modulebase/resource.go b/pkg/mcclient/modulebase/resource.go index 6f748119d8..db0b94d658 100644 --- a/pkg/mcclient/modulebase/resource.go +++ b/pkg/mcclient/modulebase/resource.go @@ -347,8 +347,15 @@ func (this *ResourceManager) BatchCreateInContexts(session *mcclient.ClientSessi ret := make([]SubmitResult, count) respbody, err := this._post(session, path, body, this.KeywordPlural) if err != nil { - for i := 0; i < count; i++ { - ret[i] = SubmitResult{Status: 500, Data: jsonutils.NewString(err.Error())} + jsonErr, ok := err.(*httputils.JSONClientError) + if ok { + for i := 0; i < count; i++ { + ret[i] = SubmitResult{Status: jsonErr.Code, Data: jsonutils.Marshal(jsonErr)} + } + } else { + for i := 0; i < count; i++ { + ret[i] = SubmitResult{Status: 500, Data: jsonutils.NewString(err.Error())} + } } return ret }