scheduler: use compute models

This commit is contained in:
Zexi Li
2019-03-20 12:03:29 +08:00
parent 6c4f925353
commit 19c341c297
65 changed files with 1849 additions and 1271 deletions
@@ -1,17 +1,13 @@
package predicates
import (
"fmt"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
computemodels "yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/plugin"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/scheduler/core"
"yunion.io/x/onecloud/pkg/scheduler/core/score"
"yunion.io/x/onecloud/pkg/scheduler/db/models"
"yunion.io/x/onecloud/pkg/util/conditionparser"
)
// NOTE: Aggregate Description
@@ -26,101 +22,16 @@ import (
type AggregatePredicate struct {
BasePredicate
plugin.BasePlugin
AggregateHosts hostsAggregatesMap
RequireAggregates []api.Aggregate
ExcludeAggregates []api.Aggregate
AvoidAggregates []api.Aggregate
PreferAggregates []api.Aggregate
AggregateMap map[string]api.Aggregate
SchedtagPredicate *SchedtagPredicate
}
type hostAggregates []*models.Aggregate
type hostsAggregatesMap map[string]hostAggregates
func (p *AggregatePredicate) Name() string {
return "host_aggregate"
}
func (p *AggregatePredicate) Clone() core.FitPredicate {
return &AggregatePredicate{
AggregateMap: make(map[string]api.Aggregate, 0),
}
}
func getHostAndServerSchedDesc(u *core.Unit, c core.Candidater) *jsonutils.JSONDict {
ret := jsonutils.NewDict()
hostSchedDesc := c.GetSchedDesc()
srvSchedDesc := jsonutils.Marshal(u.SchedData())
ret.Add(hostSchedDesc, "host")
ret.Add(srvSchedDesc, "server")
return ret
}
func getHostDynamicSchedtags(u *core.Unit, c core.Candidater) ([]*models.Aggregate, error) {
schedDesc := getHostAndServerSchedDesc(u, c)
dynamicTags, err := models.FetchEnabledDynamicschedtags()
if err != nil {
return nil, err
}
aggs := []*models.Aggregate{}
for _, tag := range dynamicTags {
matched, err := conditionparser.Eval(tag.Condition, schedDesc)
if err != nil {
log.Errorf("Condition parse eval: condition: %q, desc: %s, error: %v", tag.Condition, schedDesc, err)
continue
}
if !matched {
continue
}
aggregate, err := tag.FetchSchedTag()
if err != nil {
log.Errorf("Get dynamic schedtag %q error: %v", tag.SchedtagId, err)
continue
}
aggs = append(aggs, aggregate)
}
return aggs, nil
}
func mergeHostSchedtags(c core.Candidater, staticTags, dynamicTags []*models.Aggregate) []*models.Aggregate {
isIn := func(tags []*models.Aggregate, dt *models.Aggregate) bool {
for _, t := range tags {
if t.ID == dt.ID {
return true
}
}
return false
}
ret := []*models.Aggregate{}
ret = append(ret, staticTags...)
for _, dt := range dynamicTags {
if !isIn(staticTags, dt) {
ret = append(ret, dt)
log.Debugf("Append dynamic schedtag %s to host %q", dt, c.IndexKey())
}
}
return ret
}
func hostsAggregatesInfo(u *core.Unit, cs []core.Candidater) (hostsAggregatesMap, []*models.Aggregate) {
ret := make(map[string]hostAggregates, 0)
allAggs := make([]*models.Aggregate, 0)
for _, c := range cs {
hostAggs := c.GetHostAggregates()
dynamicMatchedAggs, err := getHostDynamicSchedtags(u, c)
if err != nil {
log.Errorf("Get host %q dynamic schedtag error: %v", c.IndexKey(), err)
} else {
hostAggs = mergeHostSchedtags(c, hostAggs, dynamicMatchedAggs)
}
ret[c.IndexKey()] = hostAggs
}
if len(cs) > 0 {
allAggs = cs[0].GetAggregates()
}
return ret, allAggs
return &AggregatePredicate{}
}
func (p *AggregatePredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
@@ -130,72 +41,16 @@ func (p *AggregatePredicate) PreExecute(u *core.Unit, cs []core.Candidater) (boo
return false, nil
}
hsMap, allAggs := hostsAggregatesInfo(u, cs)
p.AggregateHosts = hsMap
appendedAggIds := make(map[string]int, len(data.Aggregates))
for _, aggregate := range data.Aggregates {
switch aggregate.Strategy {
case api.AggregateStrategyRequire:
p.RequireAggregates = append(p.RequireAggregates, aggregate)
case api.AggregateStrategyExclude:
p.ExcludeAggregates = append(p.ExcludeAggregates, aggregate)
case api.AggregateStrategyPrefer:
p.PreferAggregates = append(p.PreferAggregates, aggregate)
case api.AggregateStrategyAvoid:
p.AvoidAggregates = append(p.AvoidAggregates, aggregate)
}
p.AggregateMap[aggregate.Idx] = aggregate
appendedAggIds[aggregate.Idx] = 1
}
for _, aggregate := range allAggs {
_, nameOk := appendedAggIds[aggregate.Name]
_, idOk := appendedAggIds[aggregate.ID]
if !(nameOk || idOk) {
agg := api.Aggregate{Idx: aggregate.ID, Strategy: aggregate.DefaultStrategy}
switch agg.Strategy {
case api.AggregateStrategyRequire:
p.RequireAggregates = append(p.RequireAggregates, agg)
case api.AggregateStrategyExclude:
p.ExcludeAggregates = append(p.ExcludeAggregates, agg)
case api.AggregateStrategyPrefer:
p.PreferAggregates = append(p.PreferAggregates, agg)
case api.AggregateStrategyAvoid:
p.AvoidAggregates = append(p.AvoidAggregates, agg)
}
}
allAggs, err := GetAllSchedtags(computemodels.HostManager.KeywordPlural())
if err != nil {
return false, err
}
p.SchedtagPredicate = NewSchedtagPredicate(data.Schedtags, allAggs)
u.AppendSelectPlugin(p)
return true, nil
}
func getHostAggregateCount(inAggs []api.Aggregate, hAggs []*models.Aggregate, strategy string) (countMap map[string]int) {
countMap = make(map[string]int)
in := func(hAgg *models.Aggregate, inAggs []api.Aggregate) bool {
for _, agg := range inAggs {
if agg.Idx == hAgg.ID || agg.Idx == hAgg.Name {
return true
}
}
return false
}
for _, hAgg := range hAggs {
if in(hAgg, inAggs) {
countMap[fmt.Sprintf("%s:%s:%s", hAgg.ID, hAgg.Name, strategy)]++
}
}
return
}
func (p *AggregatePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := NewPredicateHelper(p, u, c)
@@ -206,66 +61,46 @@ func (p *AggregatePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []c
return h.GetResult()
}
type schedtagCandidateW struct {
core.Candidater
schedData *api.SchedData
}
func (w schedtagCandidateW) GetDynamicSchedDesc() *jsonutils.JSONDict {
ret := jsonutils.NewDict()
hostSchedDesc := w.GetSchedDesc()
srvSchedDesc := jsonutils.Marshal(w.schedData)
ret.Add(hostSchedDesc, "host")
ret.Add(srvSchedDesc, "server")
return ret
}
func (w schedtagCandidateW) GetSchedtags() []computemodels.SSchedtag {
return w.Getter().HostSchedtags()
}
func (w schedtagCandidateW) ResourceType() string {
return computemodels.HostManager.KeywordPlural()
}
func (p *AggregatePredicate) exec(h *PredicateHelper) string {
ahs := p.AggregateHosts
candidateID := h.Candidate.IndexKey()
log.V(10).Debugf(">>>> ExcludeAggregates: %#v, RequireAggregates: %#v, AvoidAggregates: %#v, PreferAggregates: %#v, candidateID: %v", p.ExcludeAggregates, p.RequireAggregates, p.AvoidAggregates, p.PreferAggregates, candidateID)
if len(p.ExcludeAggregates) > 0 {
inExclude := func(a *models.Aggregate) bool {
for _, agg := range p.ExcludeAggregates {
if agg.Idx == a.ID || agg.Idx == a.Name {
return true
}
}
return false
}
if ah, ok := ahs[candidateID]; ok {
for _, a := range ah {
if inExclude(a) {
return fmt.Sprintf("exclude by aggregate: '%s:%s'", a.Name, a.ID)
}
}
}
}
if len(p.RequireAggregates) > 0 {
var as []*models.Aggregate = nil
if ah, ok := ahs[candidateID]; ok {
as = ah
}
inRequire := func(agg api.Aggregate) bool {
for _, a := range as {
if a.ID == agg.Idx || a.Name == agg.Idx {
return true
}
}
return false
}
for _, agg := range p.RequireAggregates {
if !inRequire(agg) {
return fmt.Sprintf("need aggregate: '%s'", agg.Idx)
}
}
if err := p.SchedtagPredicate.Check(
schedtagCandidateW{
Candidater: h.Candidate,
schedData: h.Unit.SchedData(),
},
); err != nil {
return err.Error()
}
return ""
}
func (p *AggregatePredicate) OnPriorityEnd(u *core.Unit, c core.Candidater) {
hostAggs, ok := p.AggregateHosts[c.IndexKey()]
if !ok {
return
}
hostAggs := c.Getter().HostSchedtags()
avoidCountMap := getHostAggregateCount(p.AvoidAggregates, hostAggs, api.AggregateStrategyAvoid)
preferCountMap := getHostAggregateCount(p.PreferAggregates, hostAggs, api.AggregateStrategyPrefer)
avoidCountMap := GetSchedtagCount(p.SchedtagPredicate.GetAvoidTags(), hostAggs, api.AggregateStrategyAvoid)
preferCountMap := GetSchedtagCount(p.SchedtagPredicate.GetPreferTags(), hostAggs, api.AggregateStrategyPrefer)
setScore := func(aggCountMap map[string]int, postiveScore bool) {
stepScore := core.PriorityStep
@@ -1,9 +1,9 @@
package baremetal
import (
o "yunion.io/x/onecloud/cmd/scheduler/options"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates"
"yunion.io/x/onecloud/pkg/scheduler/core"
o "yunion.io/x/onecloud/pkg/scheduler/options"
)
type BasePredicate struct {
@@ -61,29 +61,29 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
var errMsgs []string
for _, network := range candidate.Networks {
appendError := func(errMsg string) {
errMsgs = append(errMsgs, fmt.Sprintf("%s: %s", network.ID, errMsg))
errMsgs = append(errMsgs, fmt.Sprintf("%s: %s", network.Id, errMsg))
}
if !((network.Ports > 0 || isMigrate()) && network.IsExit == exit) {
if !((network.GetPorts() > 0 || isMigrate()) && network.IsExitNetwork() == exit) {
appendError(predicates.ErrNoPorts)
}
if wire != "" && !utils.HasPrefix(wire, network.Wire) && !utils.HasPrefix(wire, network.WireID) { // re
if wire != "" && !utils.HasPrefix(wire, network.WireId) && !utils.HasPrefix(wire, network.GetWire().GetName()) { // re
appendError(predicates.ErrWireIsNotMatch)
}
if (!private && network.IsPublic) || (private && !network.IsPublic && network.TenantID == schedData.OwnerTenantID) {
if (!private && network.IsPublic) || (private && !network.IsPublic && network.ProjectId == schedData.OwnerTenantID) {
// TODO: support reservedNetworks
reservedNetworks := 0
restPort := int64(network.Ports - reservedNetworks)
restPort := int64(network.GetPorts() - reservedNetworks)
if restPort == 0 {
appendError("not enough network port")
continue
}
counter := u.CounterManager.GetOrCreate("net:"+network.ID, func() core.Counter {
counter := u.CounterManager.GetOrCreate("net:"+network.Id, func() core.Counter {
return core.NewNormalCounter(restPort)
})
u.SharedResourceManager.Add(network.ID, counter)
u.SharedResourceManager.Add(network.Id, counter)
counters.Add(counter)
p.SelectedNetworks.Store(network.ID, counter.GetCount())
p.SelectedNetworks.Store(network.Id, counter.GetCount())
return ""
} else {
appendError(predicates.ErrNotOwner)
@@ -105,7 +105,7 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
return isRandomNetworkAvailable(network.Private, network.Exit, network.Wire)
}
for _, net := range candidate.Networks {
if (network.Idx == net.ID || network.Idx == net.Name) && (net.IsPublic || net.TenantID == schedData.OwnerTenantID) && (net.Ports > 0 || isMigrate()) {
if (network.Idx == net.Id || network.Idx == net.Name) && (net.IsPublic || net.ProjectId == schedData.OwnerTenantID) && (net.GetPorts() > 0 || isMigrate()) {
h.SetCapacity(1)
return ""
}
@@ -0,0 +1,168 @@
package predicates
import (
"fmt"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/util/errors"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/plugin"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/scheduler/core"
)
type DiskStoragesMap map[int][]*api.CandidateStorage
type CandidateDiskStoragesMap map[string]DiskStoragesMap
type DiskSchedtagPredicate struct {
BasePredicate
plugin.BasePlugin
SchedtagPredicate *SchedtagPredicate
CandidateDiskStoragesMap CandidateDiskStoragesMap
}
func (p *DiskSchedtagPredicate) Name() string {
return "disk_schedtag"
}
func (p *DiskSchedtagPredicate) Clone() core.FitPredicate {
return &DiskSchedtagPredicate{
CandidateDiskStoragesMap: make(map[string]DiskStoragesMap),
}
}
func (p *DiskSchedtagPredicate) getSchedtagDisks(disks []*api.Disk) ([]*api.Disk, []*api.Disk) {
noTagDisk := make([]*api.Disk, 0)
tagDisk := make([]*api.Disk, 0)
for _, d := range disks {
if len(d.Schedtags) != 0 {
tagDisk = append(tagDisk, d)
} else {
noTagDisk = append(noTagDisk, d)
}
}
return noTagDisk, tagDisk
}
func (p *DiskSchedtagPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
disks := u.SchedData().Disks
if len(disks) == 0 {
return false, nil
}
// always select each storages to disks
u.AppendSelectPlugin(p)
return true, nil
}
type schedtagStorageW struct {
candidater *api.CandidateStorage
disk *api.Disk
}
func (w schedtagStorageW) IndexKey() string {
return fmt.Sprintf("%d:%s", w.disk.Size, w.disk.Backend)
}
func (w schedtagStorageW) GetDynamicSchedDesc() *jsonutils.JSONDict {
return nil
}
func (w schedtagStorageW) GetSchedtags() []models.SSchedtag {
return w.candidater.Schedtags
}
func (w schedtagStorageW) ResourceType() string {
return models.StorageManager.KeywordPlural()
}
func (p *DiskSchedtagPredicate) check(d *api.Disk, s *api.CandidateStorage) (bool, error) {
allTags, err := GetAllSchedtags(models.StorageManager.KeywordPlural())
if err != nil {
return false, err
}
tagPredicate := NewSchedtagPredicate(d.Schedtags, allTags)
if err := tagPredicate.Check(
schedtagStorageW{
candidater: s,
disk: d,
},
); err != nil {
return false, err
}
return true, nil
}
func (p *DiskSchedtagPredicate) checkStorages(d *api.Disk, storages []*api.CandidateStorage) ([]*api.CandidateStorage, error) {
errs := make([]error, 0)
ret := make([]*api.CandidateStorage, 0)
for _, s := range storages {
_, err := p.check(d, s)
if err != nil {
// append err, storage not suit disk
errs = append(errs, err)
continue
}
ret = append(ret, s)
}
if len(ret) == 0 {
return nil, errors.NewAggregate(errs)
}
return ret, nil
}
func (p *DiskSchedtagPredicate) GetDiskStoragesMap(candidateId string) DiskStoragesMap {
ret, ok := p.CandidateDiskStoragesMap[candidateId]
if !ok {
ret = make(map[int][]*api.CandidateStorage)
p.CandidateDiskStoragesMap[candidateId] = ret
}
return ret
}
func (p *DiskSchedtagPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := NewPredicateHelper(p, u, c)
//noTagDisks, tagDisks := p.getSchedtagDisks(u.SchedData().Disks)
storages := c.Getter().Storages()
ds := p.GetDiskStoragesMap(c.IndexKey())
disks := u.SchedData().Disks
for _, d := range disks {
matchedStorages, err := p.checkStorages(d, storages)
if err != nil {
h.Exclude(err.Error())
}
ds[d.Index] = matchedStorages
}
return h.GetResult()
}
func (p *DiskSchedtagPredicate) OnSelectEnd(u *core.Unit, c core.Candidater, count int64) {
res := u.GetAllocatedResource(c.IndexKey())
diskStorages := p.GetDiskStoragesMap(c.IndexKey())
res.Disks = make([]*core.DiskAllocatedResource, len(diskStorages))
disks := u.SchedData().Disks
for idx, ds := range diskStorages {
res.Disks[idx] = p.allocatedDiskResource(c, disks[idx], ds)
}
log.Errorf("============OnSelectEnd %s called: %#v", c.Getter().Name(), jsonutils.Marshal(res.Disks).String())
}
func (p *DiskSchedtagPredicate) allocatedDiskResource(c core.Candidater, disk *api.Disk, storages []*api.CandidateStorage) *core.DiskAllocatedResource {
storage := p.selectStorage(disk, storages)
return &core.DiskAllocatedResource{
Index: disk.Index,
StorageId: storage.Id,
}
}
func (p *DiskSchedtagPredicate) selectStorage(d *api.Disk, storages []*api.CandidateStorage) *api.CandidateStorage {
return storages[0]
}
@@ -27,7 +27,7 @@ func (f *HypervisorPredicate) Clone() core.FitPredicate {
}
func hostHasContainerTag(c core.Candidater) bool {
aggs := c.GetHostAggregates()
aggs := c.Getter().HostSchedtags()
for _, agg := range aggs {
if agg.Name == CONTAINER_ALLOWED_TAG {
return true
@@ -8,11 +8,12 @@ import (
"yunion.io/x/pkg/util/sets"
"yunion.io/x/pkg/utils"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/plugin"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/scheduler/core"
networks "yunion.io/x/onecloud/pkg/scheduler/db/models"
)
// NetworkPredicate will filter the current network information with
@@ -56,16 +57,16 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
}
// ServerType's value is 'guest', 'container' or ''(support all type) will return true.
isMatchServerType := func(network *networks.NetworkSchedResult) bool {
isMatchServerType := func(network *models.SNetwork) bool {
return sets.NewString("guest", "", "container").Has(network.ServerType)
}
counterOfNetwork := func(u *core.Unit, n *networks.NetworkSchedResult, r int) core.Counter {
counter := u.CounterManager.GetOrCreate("net:"+n.ID, func() core.Counter {
return core.NewNormalCounter(int64(n.Ports - r))
counterOfNetwork := func(u *core.Unit, n *models.SNetwork, r int) core.Counter {
counter := u.CounterManager.GetOrCreate("net:"+n.Id, func() core.Counter {
return core.NewNormalCounter(int64(n.GetPorts() - r))
})
u.SharedResourceManager.Add(n.ID, counter)
u.SharedResourceManager.Add(n.GetId(), counter)
return counter
}
@@ -81,36 +82,36 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
errMsgs = append(errMsgs, errMsg)
}
if !isMatchServerType(n) {
if !isMatchServerType(&n) {
appendError(predicates.ErrServerTypeIsNotMatch)
}
if n.IsExit != exit {
if n.IsExitNetwork() != exit {
appendError(predicates.ErrExitIsNotMatch)
}
if !(n.Ports > 0 || isMigrate()) {
if !(n.GetPorts() > 0 || isMigrate()) {
appendError(predicates.ErrNoPorts)
}
if wire != "" && !utils.HasPrefix(wire, n.Wire) && !utils.HasPrefix(wire, n.WireID) { // re
if wire != "" && !utils.HasPrefix(wire, n.WireId) && !utils.HasPrefix(wire, n.GetWire().GetName()) { // re
appendError(predicates.ErrWireIsNotMatch)
}
if !((!private && n.IsPublic) || (private && !n.IsPublic && n.TenantID == d.OwnerTenantID)) {
if !((!private && n.IsPublic) || (private && !n.IsPublic && n.ProjectId == d.OwnerTenantID)) {
appendError(predicates.ErrNotOwner)
}
if len(errMsgs) == 0 {
// add resource
reservedNetworks := 0
counter := counterOfNetwork(u, n, reservedNetworks)
p.SelectedNetworks.Store(n.ID, counter.GetCount())
counter := counterOfNetwork(u, &n, reservedNetworks)
p.SelectedNetworks.Store(n.GetId(), counter.GetCount())
counters.Add(counter)
found = true
} else {
fullErrMsgs = append(fullErrMsgs,
fmt.Sprintf("%s: %s", n.ID, strings.Join(errMsgs, ",")),
fmt.Sprintf("%s: %s", n.Id, strings.Join(errMsgs, ",")),
)
}
}
@@ -131,7 +132,7 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
}
isNetworkAvaliable := func(n *api.Network, counters *core.MinCounters,
networks []*networks.NetworkSchedResult) string {
networks []models.SNetwork) string {
if n.Idx == "" {
counters0 := core.NewCounters()
ret_msg := isRandomNetworkAvailable(n.Private, n.Exit, n.Wire, counters0)
@@ -149,21 +150,21 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): server type not matched", net.Name, net.ID))
continue
}*/
if !(n.Idx == net.ID || n.Idx == net.Name) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): id/name not matched", net.Name, net.ID))
} else if !(net.IsPublic || net.TenantID == d.OwnerTenantID) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): not owner (%v != %v)", net.Name, net.ID, net.TenantID, d.OwnerTenantID))
} else if !(net.Ports > 0 || isMigrate()) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): ports use up", net.Name, net.ID))
if !(n.Idx == net.GetId() || n.Idx == net.GetName()) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): id/name not matched", net.Name, net.Id))
} else if !(net.IsPublic || net.ProjectId == d.OwnerTenantID) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): not owner (%v != %v)", net.Name, net.Id, net.ProjectId, d.OwnerTenantID))
} else if !(net.GetPorts() > 0 || isMigrate()) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): ports use up", net.Name, net.Id))
} else {
// add resource
reservedNetworks := 0
counter := counterOfNetwork(u, net, reservedNetworks)
counter := counterOfNetwork(u, &net, reservedNetworks)
if counter.GetCount() < d.Count {
errMsgs = append(errMsgs, fmt.Sprintf("%s: ports not enough, free: %d, required: %d", net.Name, counter.GetCount(), d.Count))
continue
}
p.SelectedNetworks.Store(net.ID, counter.GetCount())
p.SelectedNetworks.Store(net.Id, counter.GetCount())
counters.Add(counter)
return ""
}
@@ -4,6 +4,7 @@ import (
"fmt"
"strings"
"yunion.io/x/pkg/tristate"
"yunion.io/x/pkg/utils"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates"
@@ -52,7 +53,7 @@ func (p *StoragePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
isStorageAccessible := func(storage string) bool {
for _, s := range hc.Storages {
if storage == s.ID || storage == s.Name {
if storage == s.Id || storage == s.Name {
return true
}
}
@@ -82,9 +83,11 @@ func (p *StoragePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
ss := []string{}
for _, s := range hc.Storages {
if s.StorageType == backend {
total := int64(float64(s.Capacity) * s.Cmtbound)
free := total - s.UsedCapacity - s.WasteCapacity
ss = append(ss, fmt.Sprintf("(%v-%v-%v=%v)", total, s.UsedCapacity, s.WasteCapacity, free))
total := int64(float32(s.Capacity) * s.Cmtbound)
used := s.GetUsedCapacity(tristate.True)
waste := s.GetUsedCapacity(tristate.False)
free := total - int64(used) - int64(waste)
ss = append(ss, fmt.Sprintf("(%v-%v-%v=%v)", total, used, waste, free))
}
}
return strings.Join(ss, " + ")
@@ -7,9 +7,11 @@ import (
"k8s.io/client-go/kubernetes"
"yunion.io/x/pkg/util/errors"
"yunion.io/x/pkg/util/netutils"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/scheduler/cache/candidate"
"yunion.io/x/onecloud/pkg/scheduler/db/models"
)
const (
@@ -56,13 +58,13 @@ func (p *NetworkPredicate) Execute(cli *kubernetes.Clientset, pod *v1.Pod, node
return true, nil
}
func (p NetworkPredicate) checkByNetworks(nets []*models.NetworkSchedResult) error {
func (p NetworkPredicate) checkByNetworks(nets []models.SNetwork) error {
if len(nets) == 0 {
return fmt.Errorf("Network is empty")
}
errs := make([]error, 0)
for _, net := range nets {
err := p.checkByNetwork(net)
err := p.checkByNetwork(&net)
if err == nil {
return nil
}
@@ -71,23 +73,23 @@ func (p NetworkPredicate) checkByNetworks(nets []*models.NetworkSchedResult) err
return errors.NewAggregate(errs)
}
func (p NetworkPredicate) checkByNetwork(net *models.NetworkSchedResult) error {
if net.Ports <= 0 {
func (p NetworkPredicate) checkByNetwork(net *models.SNetwork) error {
if net.GetPorts() <= 0 {
return fmt.Errorf("Network %s no free IPs", net.Name)
}
if !(p.network == net.Name || p.network == net.ID) {
return fmt.Errorf("Network %s:%s or id not match %s", net.Name, net.ID, p.network)
if !(p.network == net.Name || p.network == net.Id) {
return fmt.Errorf("Network %s:%s or id not match %s", net.Name, net.Id, p.network)
}
return nil
}
func (p NetworkPredicate) checkNetworksIP(ip string, nets []*models.NetworkSchedResult) error {
func (p NetworkPredicate) checkNetworksIP(ip string, nets []models.SNetwork) error {
if len(nets) == 0 {
return fmt.Errorf("Network is empty")
}
errs := make([]error, 0)
for _, net := range nets {
err := p.checkNetworkIP(ip, net)
err := p.checkNetworkIP(ip, &net)
if err == nil {
return nil
}
@@ -96,10 +98,12 @@ func (p NetworkPredicate) checkNetworksIP(ip string, nets []*models.NetworkSched
return errors.NewAggregate(errs)
}
func (p NetworkPredicate) checkNetworkIP(ip string, net *models.NetworkSchedResult) error {
if ok, err := net.ContainsIp(ip); err != nil {
func (p NetworkPredicate) checkNetworkIP(ip string, net *models.SNetwork) error {
ipAddr, err := netutils.NewIPV4Addr(ip)
if err != nil {
return err
} else if !ok {
}
if ok := net.GetIPRange().Contains(ipAddr); !ok {
return fmt.Errorf("Network %s not contains ip %s", net.Name, ip)
}
return nil
@@ -0,0 +1,263 @@
package predicates
import (
"fmt"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/util/conditionparser"
)
type ISchedtagPredicate interface {
GetExcludeTags() []api.Schedtag
GetRequireTags() []api.Schedtag
GetAvoidTags() []api.Schedtag
GetPreferTags() []api.Schedtag
}
type ISchedtagCandidate interface {
IndexKey() string
ResourceType() string
// GetSchedtags return schedtags bind to this candidate
GetSchedtags() []models.SSchedtag
// GetDynamicSchedDesc return schedule description used by dynamic schedtags condition eval
GetDynamicSchedDesc() *jsonutils.JSONDict
}
type SchedtagPredicate struct {
requireTags []api.Schedtag
execludeTags []api.Schedtag
preferTags []api.Schedtag
avoidTags []api.Schedtag
checker *SchedtagChecker
}
func NewSchedtagPredicate(reqTags []api.Schedtag, allTags []models.SSchedtag) *SchedtagPredicate {
p := new(SchedtagPredicate)
requireTags, execludeTags, preferTags, avoidTags := GetRequestSchedtags(reqTags, allTags)
p.requireTags = requireTags
p.execludeTags = execludeTags
p.preferTags = preferTags
p.avoidTags = avoidTags
p.checker = new(SchedtagChecker)
return p
}
func (p *SchedtagPredicate) GetExcludeTags() []api.Schedtag {
return p.execludeTags
}
func (p *SchedtagPredicate) GetRequireTags() []api.Schedtag {
return p.requireTags
}
func (p *SchedtagPredicate) GetAvoidTags() []api.Schedtag {
return p.avoidTags
}
func (p *SchedtagPredicate) GetPreferTags() []api.Schedtag {
return p.preferTags
}
func (p *SchedtagPredicate) Check(candidate ISchedtagCandidate) error {
return p.checker.Check(p, candidate)
}
func GetSchedtagCount(inTags []api.Schedtag, objTags []models.SSchedtag, strategy string) (countMap map[string]int) {
countMap = make(map[string]int)
in := func(objTag models.SSchedtag, inTags []api.Schedtag) bool {
for _, tag := range inTags {
if tag.Idx == objTag.Id || tag.Idx == objTag.Name {
return true
}
}
return false
}
for _, objTag := range objTags {
if in(objTag, inTags) {
countMap[fmt.Sprintf("%s:%s:%s", objTag.Id, objTag.Name, strategy)]++
}
}
return
}
func GetAllSchedtags(resType string) ([]models.SSchedtag, error) {
tags, err := models.SchedtagManager.GetResourceSchedtags(resType)
if err != nil {
return nil, err
}
return tags, nil
}
func GetRequestSchedtags(reqTags []api.Schedtag, allTags []models.SSchedtag) (requireTags, execludeTags, preferTags, avoidTags []api.Schedtag) {
requireTags = make([]api.Schedtag, 0)
execludeTags = make([]api.Schedtag, 0)
preferTags = make([]api.Schedtag, 0)
avoidTags = make([]api.Schedtag, 0)
appendedTagIds := make(map[string]int)
appendTagByStrategy := func(tag api.Schedtag) {
switch tag.Strategy {
case models.STRATEGY_REQUIRE:
requireTags = append(requireTags, tag)
case models.STRATEGY_EXCLUDE:
execludeTags = append(execludeTags, tag)
case models.STRATEGY_PREFER:
preferTags = append(preferTags, tag)
case models.STRATEGY_AVOID:
avoidTags = append(avoidTags, tag)
}
}
for _, tag := range reqTags {
appendTagByStrategy(tag)
appendedTagIds[tag.Idx] = 1
}
for _, tag := range allTags {
_, nameOk := appendedTagIds[tag.Name]
_, idOk := appendedTagIds[tag.Id]
if !(nameOk || idOk) {
apiTag := api.Schedtag{Idx: tag.Id, Strategy: tag.DefaultStrategy}
appendTagByStrategy(apiTag)
}
}
return
}
type SchedtagChecker struct{}
type apiTags []api.Schedtag
func (t apiTags) contains(objTag models.SSchedtag) bool {
for _, tag := range t {
if tag.Idx == objTag.Id || tag.Idx == objTag.Name {
return true
}
}
return false
}
type objTags []models.SSchedtag
func (t objTags) contains(atag api.Schedtag) bool {
for _, tag := range t {
if tag.Id == atag.Idx || tag.Name == atag.Idx {
return true
}
}
return false
}
func (c *SchedtagChecker) contains(tags []api.Schedtag, objTag models.SSchedtag) bool {
for _, tag := range tags {
if tag.Idx == objTag.Id || tag.Idx == objTag.Name {
return true
}
}
return false
}
func (c *SchedtagChecker) HasIntersection(tags []api.Schedtag, objTags []models.SSchedtag) (bool, *models.SSchedtag) {
var atags apiTags = tags
for _, objTag := range objTags {
if atags.contains(objTag) {
return true, &objTag
}
}
return false, nil
}
func (c *SchedtagChecker) Contains(objectTags []models.SSchedtag, tags []api.Schedtag) (bool, *api.Schedtag) {
var otags objTags = objectTags
for _, tag := range tags {
if !otags.contains(tag) {
return false, &tag
}
}
return true, nil
}
func (p *SchedtagChecker) getDynamicSchedtags(schedDesc *jsonutils.JSONDict) ([]models.SSchedtag, error) {
if schedDesc == nil {
return []models.SSchedtag{}, nil
}
dynamicTags := models.DynamicschedtagManager.GetAllEnabledDynamicSchedtags()
tags := []models.SSchedtag{}
for _, tag := range dynamicTags {
matched, err := conditionparser.Eval(tag.Condition, schedDesc)
if err != nil {
log.Errorf("Condition parse eval: condition: %q, desc: %s, error: %v", tag.Condition, schedDesc, err)
continue
}
if !matched {
continue
}
objTag := tag.GetSchedtag()
if objTag != nil {
tags = append(tags, *objTag)
}
}
return tags, nil
}
func (c *SchedtagChecker) mergeSchedtags(candiate ISchedtagCandidate, staticTags, dynamicTags []models.SSchedtag) []models.SSchedtag {
isIn := func(tags []models.SSchedtag, dt models.SSchedtag) bool {
for _, t := range tags {
if t.Id == dt.Id {
return true
}
}
return false
}
ret := []models.SSchedtag{}
ret = append(ret, staticTags...)
for _, dt := range dynamicTags {
if !isIn(staticTags, dt) {
ret = append(ret, dt)
log.Debugf("Append dynamic schedtag %s to %s %q", dt, candiate.ResourceType(), candiate.IndexKey())
}
}
return ret
}
func (c *SchedtagChecker) GetCandidateSchedtags(candidate ISchedtagCandidate) ([]models.SSchedtag, error) {
staticTags := candidate.GetSchedtags()
dynamicTags, err := c.getDynamicSchedtags(candidate.GetDynamicSchedDesc())
if err != nil {
return nil, err
}
return c.mergeSchedtags(candidate, staticTags, dynamicTags), nil
}
func (c *SchedtagChecker) Check(p ISchedtagPredicate, candidate ISchedtagCandidate) error {
candidateTags, err := c.GetCandidateSchedtags(candidate)
if err != nil {
return err
}
execludeTags := p.GetExcludeTags()
if len(execludeTags) > 0 {
if ok, tag := c.HasIntersection(execludeTags, candidateTags); ok {
return fmt.Errorf("Execlude by schedtag: '%s:%s'", tag.Name, tag.Id)
}
}
requireTags := p.GetRequireTags()
if len(requireTags) > 0 {
if ok, tag := c.Contains(candidateTags, requireTags); !ok {
return fmt.Errorf("Need schedtag: '%s'", tag.Idx)
}
}
return nil
}
@@ -1,55 +0,0 @@
package guest
import (
"yunion.io/x/onecloud/pkg/scheduler/algorithm/priorities"
"yunion.io/x/onecloud/pkg/scheduler/core"
)
type AvoidSameClusterPriority struct {
priorities.BasePriority
ClusterTbl map[string]int
}
func (p *AvoidSameClusterPriority) Name() string {
return "guest_avoid_same_cluster"
}
func (p *AvoidSameClusterPriority) Clone() core.Priority {
return &AvoidSameClusterPriority{ClusterTbl: make(map[string]int)}
}
func (p *AvoidSameClusterPriority) PreExecute(u *core.Unit, cs []core.Candidater) (bool, []core.PredicateFailureReason, error) {
d := u.SchedData()
clusterTbl := make(map[string]int, 0)
ownerTenantID := d.OwnerTenantID
for _, c := range cs {
hc, err := p.HostCandidate(c)
if err != nil {
return false, nil, err
}
if count, ok := hc.Tenants[ownerTenantID]; ok && count > 0 {
clusterId := hc.ClusterID
if count0, ok := clusterTbl[clusterId]; ok {
clusterTbl[clusterId] = count0 + int(count)
} else {
clusterTbl[clusterId] = int(count)
}
}
}
p.ClusterTbl = clusterTbl
return true, nil, nil
}
func (p *AvoidSameClusterPriority) Map(u *core.Unit, c core.Candidater) (core.HostPriority, error) {
h := priorities.NewPriorityHelper(p, u, c)
if count, ok := p.ClusterTbl[c.Get("ClusterID").(string)]; ok {
h.SetScore(-20 * count)
}
return h.GetResult()
}