mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-21 06:09:39 +08:00
改进:增加schedpolicy和dynamicschedtag支持
This commit is contained in:
@@ -0,0 +1,176 @@
|
||||
package models
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
|
||||
"database/sql"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db"
|
||||
"yunion.io/x/onecloud/pkg/httperrors"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
"yunion.io/x/onecloud/pkg/util/conditionparser"
|
||||
)
|
||||
|
||||
type SDynamicschedtagManager struct {
|
||||
db.SStandaloneResourceBaseManager
|
||||
}
|
||||
|
||||
var DynamicschedtagManager *SDynamicschedtagManager
|
||||
|
||||
func init() {
|
||||
DynamicschedtagManager = &SDynamicschedtagManager{SStandaloneResourceBaseManager: db.NewStandaloneResourceBaseManager(SDynamicschedtag{}, "dynamicschedtags_tbl", "dynamicschedtag", "dynamicschedtags")}
|
||||
}
|
||||
|
||||
// dynamic schedtag is called before scan host candidates, dynamically adding additional schedtag to hosts
|
||||
// condition examples:
|
||||
// host.sys_load > 1.5 || host.mem_used_percent > 0.7 => "high_load"
|
||||
//
|
||||
type SDynamicschedtag struct {
|
||||
db.SStandaloneResourceBase
|
||||
|
||||
Condition string `width:"256" charset:"ascii" nullable:"false" list:"user" create:"required"`
|
||||
SchedtagId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required"`
|
||||
|
||||
Enabled bool `nullable:"false" default:"true" create:"optional" list:"user" update:"user"`
|
||||
}
|
||||
|
||||
func validateDynamicSchedtagInputData(data *jsonutils.JSONDict, create bool) error {
|
||||
condStr := jsonutils.GetAnyString(data, []string{"condition"})
|
||||
if len(condStr) == 0 && create {
|
||||
return httperrors.NewInputParameterError("empty condition")
|
||||
}
|
||||
if len(condStr) > 0 && !conditionparser.IsValid(condStr) {
|
||||
return httperrors.NewInputParameterError("invalid condition")
|
||||
}
|
||||
|
||||
schedStr := jsonutils.GetAnyString(data, []string{"schedtag", "schedtag_id"})
|
||||
if len(schedStr) == 0 && create {
|
||||
return httperrors.NewInputParameterError("missing schedtag")
|
||||
}
|
||||
if len(schedStr) > 0 {
|
||||
schedObj, err := SchedtagManager.FetchByIdOrName(nil, schedStr)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
return httperrors.NewResourceNotFoundError("schedtag %s not found", schedStr)
|
||||
} else {
|
||||
log.Errorf("fetch schedtag %s fail %s", schedStr, err)
|
||||
return httperrors.NewGeneralError(err)
|
||||
}
|
||||
}
|
||||
data.Set("schedtag_id", jsonutils.NewString(schedObj.GetId()))
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (manager *SDynamicschedtagManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerProjId string, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
|
||||
err := validateDynamicSchedtagInputData(data, true)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return manager.SStandaloneResourceBaseManager.ValidateCreateData(ctx, userCred, ownerProjId, query, data)
|
||||
}
|
||||
|
||||
func (self *SDynamicschedtag) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
|
||||
err := validateDynamicSchedtagInputData(data, false)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return self.SStandaloneResourceBase.ValidateUpdateData(ctx, userCred, query, data)
|
||||
}
|
||||
|
||||
func (self *SDynamicschedtag) getSchedtag() *SSchedtag {
|
||||
obj, err := SchedtagManager.FetchById(self.SchedtagId)
|
||||
if err != nil {
|
||||
log.Errorf("fail to fetch sched tag by id %s", err)
|
||||
return nil
|
||||
}
|
||||
return obj.(*SSchedtag)
|
||||
}
|
||||
|
||||
func (self *SDynamicschedtag) getMoreColumns(extra *jsonutils.JSONDict) *jsonutils.JSONDict {
|
||||
schedtag := self.getSchedtag()
|
||||
if schedtag != nil {
|
||||
extra.Add(jsonutils.NewString(schedtag.GetName()), "schedtag")
|
||||
}
|
||||
return extra
|
||||
}
|
||||
|
||||
func (self *SDynamicschedtag) GetCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict {
|
||||
extra := self.SStandaloneResourceBase.GetCustomizeColumns(ctx, userCred, query)
|
||||
return self.getMoreColumns(extra)
|
||||
}
|
||||
|
||||
func (self *SDynamicschedtag) GetExtraDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict {
|
||||
extra := self.SStandaloneResourceBase.GetExtraDetails(ctx, userCred, query)
|
||||
return self.getMoreColumns(extra)
|
||||
}
|
||||
|
||||
func (manager *SDynamicschedtagManager) getAllEnabledDynamicSchedtags() []SDynamicschedtag {
|
||||
rules := make([]SDynamicschedtag, 0)
|
||||
|
||||
q := DynamicschedtagManager.Query().IsTrue("enabled")
|
||||
err := db.FetchModelObjects(manager, q, &rules)
|
||||
if err != nil {
|
||||
log.Errorf("getAllEnabledDynamicSchedtags fail %s", err)
|
||||
return nil
|
||||
}
|
||||
|
||||
return rules
|
||||
}
|
||||
|
||||
func (self *SDynamicschedtag) AllowPerformEvaluate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
|
||||
return userCred.IsSystemAdmin()
|
||||
}
|
||||
|
||||
func (self *SDynamicschedtag) PerformEvaluate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
|
||||
serverStr := jsonutils.GetAnyString(data, []string{"server", "server_id", "guest", "guest_id"})
|
||||
serverObj, err := GuestManager.FetchByIdOrName(userCred, serverStr)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, httperrors.NewResourceNotFoundError("server %s not found", serverStr)
|
||||
} else {
|
||||
return nil, httperrors.NewGeneralError(err)
|
||||
}
|
||||
}
|
||||
|
||||
server := serverObj.(*SGuest)
|
||||
srvDesc := server.getSchedDesc()
|
||||
|
||||
hostStr := jsonutils.GetAnyString(data, []string{"host", "host_id"})
|
||||
hostObj, err := HostManager.FetchByIdOrName(userCred, hostStr)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, httperrors.NewResourceNotFoundError("host %s not found", serverStr)
|
||||
} else {
|
||||
return nil, httperrors.NewGeneralError(err)
|
||||
}
|
||||
}
|
||||
|
||||
host := hostObj.(*SHost)
|
||||
// TODO: to fill host scheduling information
|
||||
hostDesc := jsonutils.Marshal(host)
|
||||
|
||||
params := jsonutils.NewDict()
|
||||
params.Add(srvDesc, "server")
|
||||
params.Add(hostDesc, "host")
|
||||
|
||||
meet, err := conditionparser.Eval(self.Condition, params)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
result := jsonutils.NewDict()
|
||||
result.Add(srvDesc, "server")
|
||||
result.Add(hostDesc, "host")
|
||||
|
||||
if meet {
|
||||
result.Add(jsonutils.JSONTrue, "result")
|
||||
} else {
|
||||
result.Add(jsonutils.JSONFalse, "result")
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
@@ -681,11 +681,11 @@ func (manager *SGuestManager) ValidateCreateData(ctx context.Context, userCred m
|
||||
return nil, httperrors.NewInputParameterError("invalid aggregate_strategy")
|
||||
}
|
||||
}
|
||||
for idx := 0; data.Contains(fmt.Sprintf("srvtag.%d", idx)); idx += 1 {
|
||||
aggStr, _ := data.GetString(fmt.Sprintf("srvtag.%d", idx))
|
||||
for idx := 0; data.Contains(fmt.Sprintf("schedtag.%d", idx)); idx += 1 {
|
||||
aggStr, _ := data.GetString(fmt.Sprintf("schedtag.%d", idx))
|
||||
if len(aggStr) > 0 {
|
||||
parts := strings.Split(aggStr, ":")
|
||||
if len(parts) >= 2 && len(parts) > 0 && len(parts[1]) > 0 {
|
||||
if len(parts) >= 2 && len(parts[0]) > 0 && len(parts[1]) > 0 {
|
||||
schedtags[parts[0]] = parts[1]
|
||||
}
|
||||
}
|
||||
@@ -4567,3 +4567,35 @@ func (self *SGuest) PerformUserData(ctx context.Context, userCred mcclient.Token
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (self *SGuest) getSchedDesc() jsonutils.JSONObject {
|
||||
desc := jsonutils.NewDict()
|
||||
|
||||
desc.Add(jsonutils.NewString(self.Id), "id")
|
||||
desc.Add(jsonutils.NewString(self.Name), "name")
|
||||
desc.Add(jsonutils.NewInt(int64(self.VmemSize)), "vmem_size")
|
||||
desc.Add(jsonutils.NewInt(int64(self.VcpuCount)), "vcpu_count")
|
||||
|
||||
gds := self.GetDisks()
|
||||
if gds != nil {
|
||||
for i := 0; i < len(gds); i += 1 {
|
||||
desc.Add(jsonutils.Marshal(gds[i].ToDiskInfo()), fmt.Sprintf("disk.%d", i))
|
||||
}
|
||||
}
|
||||
|
||||
gns := self.GetNetworks()
|
||||
if gns != nil {
|
||||
for i := 0; i < len(gns); i += 1 {
|
||||
desc.Add(jsonutils.NewString(fmt.Sprintf("%s:%s", gns[i].NetworkId, gns[i].IpAddr)), fmt.Sprintf("net.%d", i))
|
||||
}
|
||||
}
|
||||
|
||||
if len(self.HostId) > 0 && regutils.MatchUUID(self.HostId) {
|
||||
desc.Add(jsonutils.NewString(self.HostId), "host_id")
|
||||
}
|
||||
|
||||
desc.Add(jsonutils.NewString(self.ProjectId), "owner_tenant_id")
|
||||
desc.Add(jsonutils.NewString(self.GetHypervisor()), "hypervisor")
|
||||
|
||||
return desc
|
||||
}
|
||||
|
||||
@@ -0,0 +1,189 @@
|
||||
package models
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
|
||||
"database/sql"
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db"
|
||||
"yunion.io/x/onecloud/pkg/httperrors"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
"yunion.io/x/onecloud/pkg/util/conditionparser"
|
||||
"yunion.io/x/pkg/utils"
|
||||
)
|
||||
|
||||
type SSchedpolicyManager struct {
|
||||
db.SStandaloneResourceBaseManager
|
||||
SInfrastructureManager
|
||||
}
|
||||
|
||||
var SchedpolicyManager *SSchedpolicyManager
|
||||
|
||||
func init() {
|
||||
SchedpolicyManager = &SSchedpolicyManager{SStandaloneResourceBaseManager: db.NewStandaloneResourceBaseManager(SSchedpolicy{}, "schedpolicies_tbl", "schedpolicy", "schedpolicies")}
|
||||
}
|
||||
|
||||
// sched policy is called before calling scheduler, add additional preferences for schedtags
|
||||
type SSchedpolicy struct {
|
||||
db.SStandaloneResourceBase
|
||||
SInfrastructure
|
||||
|
||||
Condition string `width:"256" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"`
|
||||
SchedtagId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"`
|
||||
Strategy string `width:"32" charset:"ascii" nullable:"false" list:"user" create:"required" update:"user"`
|
||||
|
||||
Enabled bool `nullable:"false" default:"true" create:"optional" list:"user" update:"user"`
|
||||
}
|
||||
|
||||
func validateSchedpolicyInputData(data *jsonutils.JSONDict, create bool) error {
|
||||
err := validateDynamicSchedtagInputData(data, create)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
strategyStr := jsonutils.GetAnyString(data, []string{"strategy"})
|
||||
if len(strategyStr) == 0 && create {
|
||||
return httperrors.NewInputParameterError("missing strategy")
|
||||
}
|
||||
|
||||
if len(strategyStr) > 0 && !utils.IsInStringArray(strategyStr, STRATEGY_LIST) {
|
||||
return httperrors.NewInputParameterError("invalid strategy %s", strategyStr)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (manager *SSchedpolicyManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerProjId string, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
|
||||
err := validateSchedpolicyInputData(data, true)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return manager.SStandaloneResourceBaseManager.ValidateCreateData(ctx, userCred, ownerProjId, query, data)
|
||||
}
|
||||
|
||||
func (self *SSchedpolicy) ValidateUpdateData(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
|
||||
err := validateSchedpolicyInputData(data, false)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return self.SStandaloneResourceBase.ValidateUpdateData(ctx, userCred, query, data)
|
||||
}
|
||||
|
||||
func (self *SSchedpolicy) getSchedtag() *SSchedtag {
|
||||
obj, err := SchedtagManager.FetchById(self.SchedtagId)
|
||||
if err != nil {
|
||||
log.Errorf("fail to fetch sched tag by id %s", err)
|
||||
return nil
|
||||
}
|
||||
return obj.(*SSchedtag)
|
||||
}
|
||||
|
||||
func (self *SSchedpolicy) getMoreColumns(extra *jsonutils.JSONDict) *jsonutils.JSONDict {
|
||||
schedtag := self.getSchedtag()
|
||||
if schedtag != nil {
|
||||
extra.Add(jsonutils.NewString(schedtag.GetName()), "schedtag")
|
||||
}
|
||||
return extra
|
||||
}
|
||||
|
||||
func (self *SSchedpolicy) GetCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict {
|
||||
extra := self.SStandaloneResourceBase.GetCustomizeColumns(ctx, userCred, query)
|
||||
return self.getMoreColumns(extra)
|
||||
}
|
||||
|
||||
func (self *SSchedpolicy) GetExtraDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict {
|
||||
extra := self.SStandaloneResourceBase.GetExtraDetails(ctx, userCred, query)
|
||||
return self.getMoreColumns(extra)
|
||||
}
|
||||
|
||||
func (manager *SSchedpolicyManager) getAllEnabledPolicies() []SSchedpolicy {
|
||||
policies := make([]SSchedpolicy, 0)
|
||||
|
||||
q := SchedpolicyManager.Query().IsTrue("enabled")
|
||||
err := db.FetchModelObjects(manager, q, &policies)
|
||||
if err != nil {
|
||||
log.Errorf("getAllEnabledPolicies fail %s", err)
|
||||
return nil
|
||||
}
|
||||
|
||||
return policies
|
||||
}
|
||||
|
||||
func (self *SSchedpolicy) AllowPerformEvaluate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
|
||||
return userCred.IsSystemAdmin()
|
||||
}
|
||||
|
||||
func (self *SSchedpolicy) PerformEvaluate(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
|
||||
serverStr := jsonutils.GetAnyString(data, []string{"server", "server_id", "guest", "guest_id"})
|
||||
serverObj, err := GuestManager.FetchByIdOrName(userCred, serverStr)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, httperrors.NewResourceNotFoundError("server %s not found", serverStr)
|
||||
} else {
|
||||
return nil, httperrors.NewGeneralError(err)
|
||||
}
|
||||
}
|
||||
|
||||
server := serverObj.(*SGuest)
|
||||
desc := server.getSchedDesc()
|
||||
|
||||
params := jsonutils.NewDict()
|
||||
params.Add(desc, "server")
|
||||
|
||||
meet, err := conditionparser.Eval(self.Condition, params)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
result := jsonutils.NewDict()
|
||||
result.Add(desc, "server")
|
||||
if meet {
|
||||
result.Add(jsonutils.JSONTrue, "result")
|
||||
} else {
|
||||
result.Add(jsonutils.JSONFalse, "result")
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func ApplySchedPolicies(params *jsonutils.JSONDict) *jsonutils.JSONDict {
|
||||
policies := SchedpolicyManager.getAllEnabledPolicies()
|
||||
if policies == nil {
|
||||
log.Errorf("getAllEnabledPolicies fail")
|
||||
return params
|
||||
}
|
||||
|
||||
schedtags := make(map[string]string)
|
||||
|
||||
if params.Contains("aggregate_strategy") {
|
||||
err := params.Unmarshal(&schedtags, "aggregate_strategy")
|
||||
if err != nil {
|
||||
log.Errorf("unmarshall aggregate_strategy fail %s", err)
|
||||
return params
|
||||
}
|
||||
log.Infof("original sched tag %#v", schedtags)
|
||||
}
|
||||
|
||||
input := jsonutils.NewDict()
|
||||
input.Add(params, "server")
|
||||
|
||||
for i := 0; i < len(policies); i += 1 {
|
||||
meet, err := conditionparser.Eval(policies[i].Condition, input)
|
||||
if err == nil && meet {
|
||||
st := policies[i].getSchedtag()
|
||||
if st != nil {
|
||||
schedtags[st.Name] = policies[i].Strategy
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
newSchedtags := jsonutils.Marshal(schedtags)
|
||||
log.Infof("updated sched tag %s", newSchedtags)
|
||||
|
||||
ret := jsonutils.NewDict()
|
||||
ret.Add(newSchedtags, "aggregate_strategy")
|
||||
|
||||
return ret
|
||||
}
|
||||
@@ -30,6 +30,7 @@ var STRATEGY_LIST = []string{STRATEGY_REQUIRE, STRATEGY_EXCLUDE, STRATEGY_PREFER
|
||||
|
||||
type SSchedtagManager struct {
|
||||
db.SStandaloneResourceBaseManager
|
||||
SInfrastructureManager
|
||||
}
|
||||
|
||||
var SchedtagManager *SSchedtagManager
|
||||
@@ -40,6 +41,7 @@ func init() {
|
||||
|
||||
type SSchedtag struct {
|
||||
db.SStandaloneResourceBase
|
||||
SInfrastructure
|
||||
|
||||
DefaultStrategy string `width:"16" charset:"ascii" nullable:"true" default:"" list:"user" update:"admin" create:"admin_optional"` // Column(VARCHAR(16, charset='ascii'), nullable=True, default='')
|
||||
}
|
||||
@@ -49,11 +51,7 @@ func (manager *SSchedtagManager) AllowListItems(ctx context.Context, userCred mc
|
||||
}
|
||||
|
||||
func (self *SSchedtag) AllowGetDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool {
|
||||
return userCred.IsSystemAdmin()
|
||||
}
|
||||
|
||||
func (manager *SSchedtagManager) AllowCreateItem(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
|
||||
return userCred.IsSystemAdmin()
|
||||
return true
|
||||
}
|
||||
|
||||
func (manager *SSchedtagManager) ValidateSchedtags(userCred mcclient.TokenCredential, schedtags map[string]string) (map[string]string, error) {
|
||||
@@ -113,6 +111,12 @@ func (self *SSchedtag) ValidateDeleteCondition(ctx context.Context) error {
|
||||
if self.GetHostCount() > 0 {
|
||||
return httperrors.NewNotEmptyError("Tag is associated with hosts")
|
||||
}
|
||||
if self.getDynamicSchedtagCount() > 0 {
|
||||
return httperrors.NewNotEmptyError("tag has dynamic rules")
|
||||
}
|
||||
if self.getSchedPoliciesCount() > 0 {
|
||||
return httperrors.NewNotEmptyError("tag is associate with sched policies")
|
||||
}
|
||||
return self.SStandaloneResourceBase.ValidateDeleteCondition(ctx)
|
||||
}
|
||||
|
||||
@@ -151,8 +155,18 @@ func (self *SSchedtag) GetHostCount() int {
|
||||
return HostschedtagManager.Query().Equals("schedtag_id", self.Id).Count()
|
||||
}
|
||||
|
||||
func (self *SSchedtag) getSchedPoliciesCount() int {
|
||||
return SchedpolicyManager.Query().Equals("schedtag_id", self.Id).Count()
|
||||
}
|
||||
|
||||
func (self *SSchedtag) getDynamicSchedtagCount() int {
|
||||
return DynamicschedtagManager.Query().Equals("schedtag_id", self.Id).Count()
|
||||
}
|
||||
|
||||
func (self *SSchedtag) getMoreColumns(extra *jsonutils.JSONDict) *jsonutils.JSONDict {
|
||||
extra.Add(jsonutils.NewInt(int64(self.GetHostCount())), "host_count")
|
||||
extra.Add(jsonutils.NewInt(int64(self.getDynamicSchedtagCount())), "dynamic_schedtag_count")
|
||||
extra.Add(jsonutils.NewInt(int64(self.getSchedPoliciesCount())), "schedpolicy_count")
|
||||
return extra
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user