Merge pull request #4488 from swordqiu/hotfix/qj-scheduler-quota-support

Hotfix/qj scheduler quota support
This commit is contained in:
Zexi Li
2020-01-03 17:13:22 +08:00
committed by GitHub
31 changed files with 618 additions and 145 deletions
+1 -1
View File
@@ -24,7 +24,7 @@ import (
)
type ComputeQuotaKeys struct {
RegionQuotaKeys
ZoneQuotaKeys
Hypervisor string `help:"hypervisor" choices:"kvm|baremetal"`
}
+2
View File
@@ -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 {
+49 -32
View File
@@ -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"
@@ -38,15 +39,10 @@ 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"
)
var (
CancelUsages func(ctx context.Context, userCred mcclient.TokenCredential, usages []IUsage)
)
type DBModelDispatcher struct {
modelManager IModelManager
}
@@ -1104,9 +1100,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,54 +1181,69 @@ 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 {
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)
+3 -6
View File
@@ -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
+13 -2
View File
@@ -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 ""
}
+87
View File
@@ -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
}
+2
View File
@@ -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)
}
+48 -3
View File
@@ -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)
@@ -271,3 +268,51 @@ 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()
used := manager.newQuota()
err := manager.usageStore.GetQuota(ctx, keys, used)
if err != nil {
return 0, errors.Wrap(err, "manager.usageStore.GetQuotaByKeys")
}
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
}
used.Add(pendings[i])
}
err = used.Exceed(request, quota)
if err != nil {
return 0, nil
}
quota.Sub(used)
cnt := quota.Allocable(request)
return cnt, nil
}
+14 -1
View File
@@ -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"
@@ -33,6 +34,8 @@ func init() {
quotaManagerTable = make(map[reflect.Type]IQuotaManager)
db.CancelUsages = CancelUsages
db.CancelPendingUsagesInContext = cancelPendingUsagesInContext
db.InitPendingUsagesInContext = initPendingUsagesInContext
}
func Register(manager IQuotaManager) {
@@ -59,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) {
@@ -75,3 +83,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)
}
+34
View File
@@ -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
)
+2 -2
View File
@@ -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")
+1 -1
View File
@@ -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 {
+1
View File
@@ -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")))
+15 -34
View File
@@ -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,20 @@ 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, &regionReq)
if err != nil {
return nil, nil, err
return errors.Wrap(err, "quotas.CheckSetPendingQuota")
}
return &req, &regionReq, nil
return nil
}
func (self *SGuest) checkUpdateQuota(ctx context.Context, userCred mcclient.TokenCredential, vcpuCount int, vmemSize int) (quotas.IQuota, error) {
@@ -1467,7 +1448,7 @@ func getGuestResourceRequirements(
eBw := 0
iBw := 0
for _, netConfig := range input.Networks {
if isExitNetworkInfo(netConfig) {
if IsExitNetworkInfo(netConfig) {
eNicCnt += 1
eBw += netConfig.BwLimit
} else {
+2 -2
View File
@@ -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 {
+1 -1
View File
@@ -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)
+9
View File
@@ -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 {
+25 -1
View File
@@ -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 {
+42
View File
@@ -312,6 +312,48 @@ 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
}
if self.Loadbalancer >= 0 && squota.Loadbalancer > 0 && (cnt < 0 || cnt > self.Loadbalancer/squota.Loadbalancer) {
cnt = self.Loadbalancer / squota.Loadbalancer
}
return cnt
}
func (self *SRegionQuota) Update(quota quotas.IQuota) {
squota := quota.(*SRegionQuota)
if squota.Port > 0 {
+4
View File
@@ -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 {
+3
View File
@@ -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)
@@ -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)
@@ -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)
+14 -2
View File
@@ -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)
+13 -55
View File
@@ -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(&regionUsage, 1)
schedInput.PendingUsages = []jsonutils.JSONObject{
jsonutils.Marshal(&computeUsage),
jsonutils.Marshal(&regionUsage),
}
params := jsonutils.Marshal(schedInput).(*jsonutils.JSONDict)
task.SetStage("OnScheduleComplete", params)
@@ -158,31 +139,8 @@ func doScheduleObjects(
}
func cancelPendingUsage(ctx context.Context, task IScheduleTask) {
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, 0)
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)
}
}
ClearTaskPendingUsage(ctx, task.(taskman.ITask))
ClearTaskPendingRegionUsage(ctx, task.(taskman.ITask))
}
func onSchedulerRequestFail(
+9
View File
@@ -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 {
+9 -2
View File
@@ -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
}
@@ -0,0 +1,130 @@
// 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(&regionPending)
}
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 computeCnt >= 0 && (minCnt < 0 || minCnt > computeCnt) {
minCnt = computeCnt
}
}
if !regionPending.IsEmpty() {
regionCnt, _ := quotas.GetQuotaCount(ctx, &regionQuota, regionPending.GetKeys())
if regionCnt >= 0 && (minCnt < 0 || minCnt > regionCnt) {
minCnt = regionCnt
}
}
if minCnt >= 0 {
h.SetCapacity(int64(minCnt))
}
return h.GetResult()
}
@@ -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{}),
)
}
+31
View File
@@ -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,9 @@ func (b BaseHostDesc) GetResourceType() string {
func (b *BaseHostDesc) fillCloudProvider(host *computemodels.SHost) error {
b.Cloudprovider = host.GetCloudprovider()
if b.Cloudprovider != nil {
b.Cloudaccount = b.Cloudprovider.GetCloudaccount()
}
return nil
}
@@ -417,6 +425,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 {
+2
View File
@@ -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